--- title: "05-Elasticsearch 学习笔记" created: 2026-01-09 tags: - 博客 --- # Elasticsearch 学习笔记 > Elasticsearch 写入时先将每篇文章分词,再反向建立 “单个词汇→包含该词汇的所有文章” 的倒排索引,同时对词汇排序以支撑后续高效查询;搜索时则先借助内存中的 Term Index 前缀定位 + 二分查找快速找到目标词汇对应的文章 ID 列表,再根据 AND/OR 等搜索条件对多个词汇的文章 ID 列表取交集或并集,最后依据筛选后的文章 ID 提取完整内容并返回。 ## **一、顶层设计:为什么需要 Elasticsearch?** ### **1.1 ES 的定位** ![[7-Blog/后端与微服务/assets/Elasticsearch_在系统中的定位-9b2df756.jpg]] ### **1.2 传统数据库面临的问题** ![[7-Blog/后端与微服务/assets/传统关系型数据库的瓶颈-2e271bb1.jpg]] - **模糊查询性能瓶颈**:MySQL 等关系型数据库使用 `LIKE %keyword%` 进行左模糊或全模糊查询时,无法利用索引(Index),会导致**全表扫描**,在百万级数据量下性能急剧下降。 - **复杂维度筛选困难**:当面临海量数据的多维度筛选、统计、标签聚合时,SQL 语句变得极度复杂且执行效率低下。 - **文本相关性缺失**:传统数据库只能做“精准匹配”或“简单的包含匹配”,无法计算**相关性得分(Score)**,无法按“搜索结果匹配度”对结果进行排序。 - **分词能力弱**:无法处理同义词、纠错、中文分词(如将“由于”和“游泳”区分开)等复杂的自然语言处理需求。 ### **1.3 Elasticsearch 解决的核心问题** ![[7-Blog/后端与微服务/assets/Elasticsearch_核心价值-0a0cdb5f.jpg]] - **性能全文检索**:基于**倒排索引**,实现亿级数据毫秒级响应。 - **高可用与横向扩展**:天生的分布式架构,通过增加节点即可线性提升存储容量和计算能力。 - **复杂聚合分析**:提供强大的聚合(Aggregations)框架,能够替代部分 OLAP 场景,实时统计数据分布。 - **相关性排序**:基于 TF-IDF / BM25 算法,提供符合人类直觉的搜索结果排名。 ## **二、核心原理:如何解决问题?** ### **2.1 倒排索引(Inverted Index)** 这是 ES 快如闪电的根本原因。 - **正排索引(Forward Index)**:文档 ID -> 文档内容(类似 MySQL 主键查询)。 - **倒排索引**:单词(Term)-> 包含该单词的文档 ID 列表。 - ES 在写入数据时,会通过**分词器(Analyzer)**将文本拆解为单词,建立索引。 - 查询时,直接根据单词找到文档 ID 列表,并通过位运算快速合并结果,无需扫描全表。 ![[7-Blog/后端与微服务/assets/倒排索引原理详解-5d950c3c.jpg]] ### **2.2 Term Index(词项索引)** ![[7-Blog/后端与微服务/assets/Term_Index_加速原理-fe766372.jpg]] - **面临的问题**:Term Dictionary(词项字典)数据量极大(千万/亿级),无法全部放入内存,必须存在磁盘。每次查询都去磁盘翻字典,I/O 消耗巨大。 - **解决方案**:**Term Index** 是一个基于词项前缀构建的**精简目录树**(类似 Trie 树或 FST - Finite State Transducers)。 - **内存驻留**:它体积非常小,可以完全加载到 **内存(RAM)** 中。 - **加速原理**:查询时,先在内存中查 Term Index,找到该 Term 在磁盘 Term Dictionary 中的大概位置(Offset 偏移量),然后再去磁盘读取具体内容。 - *类比*:Term Dictionary 是字典本身(在磁盘),Term Index 是字典的“拼音首字母索引页”(在内存)。 ### **2.3 存储结构** ![[7-Blog/后端与微服务/assets/ES_存储结构详解-c3bd152d.jpg]] 当倒排索引帮我们找到文档 ID 后,我们还需要获取内容或进行排序。 - **Stored Fields(行式存储)** - **用途**:用于存储**原始文档内容**(`_source` 字段)。 - **特点**:行式存储,适合展示完整信息,但不适合聚合分析(读取很多冗余数据)。 - **Doc Values(列式存储)** - **用途**:专用于**排序(Sorting)和聚合(Aggregations)**。 - **原理**:空间换时间。将散落在不同文档中的同一个字段值,集中在一起列式存储。 - **优势**:排序或计算平均值时,CPU 可以连续读取内存地址,极大提升效率,且对操作系统文件缓存(OS Cache)非常友好。 ### **2.4 Segment 与 Lucene** ![[7-Blog/后端与微服务/assets/Segment_与_Lucene_架构-6f82b8e4.jpg]] - **Lucene**:ES 的底层核心库。一个 ES 的分片(Shard)本质上就是一个完整的 **Lucene 索引**。 - **Segment(段)**: - **最小单元**:Lucene 内部由多个 Segment 组成。一个 Segment 包含了倒排索引、Term Index、Stored Fields 等所有结构,是具备完整搜索功能的最小单元。 - **不可变性(Immutable)**:Segment 一旦生成,**不可修改**。 - **写入与合并**:新增数据会生成新的 Segment;删除数据只是打上 `.del` 标记(逻辑删除)。ES 会在后台自动进行 **Segment Merging(段合并)**,将多个小 Segment 合并为大 Segment,同时物理剔除被标记删除的数据。 ## **三、分布式架构设计** ![[7-Blog/后端与微服务/assets/ES_分布式架构设计-01752588.jpg]] - **Cluster(集群)**:由多个节点组成。 - **Node(节点)**:单个服务实例。 - **Index(索引)**:逻辑上的数据集合(类比 Database)。 - **Shard(分片)**:数据被切分成多个分片存储在不同节点,实现并行计算和存储扩展。 - **Replica(副本)**:分片的备份,用于提高可用性(HA)和读取吞吐量。 ### **3.1 高性能优化:分片(Shard)** ![[7-Blog/后端与微服务/assets/高性能优化:分片机制-6bcd256e.jpg]] - **问题**:单个 Index 数据量过大(如 1TB),单机硬盘存不下,且搜索时单线程扫描太慢。 - **解决方案**:**分片机制**。 - 将一个 Index 逻辑拆分为多个 Shard(分片)。 - 每个 Shard 是一个独立的 Lucene 实例。 - **优势**:读写压力被分散到多个 Shard 上并行处理,极大提升吞吐量。 ### **3.2 高扩展性优化:多节点(Node)** ![[7-Blog/后端与微服务/assets/高扩展性优化:多节点部署-fa229e0c.jpg]] - **原理**:**横向扩展(Scale Out)**。 - **机制**:当数据量增长,分片增多,单机 CPU/内存吃紧时,可以向集群中加入新的机器(Node)。 - **自动平衡**:ES 会自动感知新节点,并将部分 Shard 迁移过去,实现负载均衡。 ### **3.3 高可用优化:副本(Replica)** ![[7-Blog/后端与微服务/assets/高可用优化:副本机制-27535f46.jpg]] - **角色区分**: - **Primary Shard(主分片)**:负责处理写入请求。 - **Replica Shard(副本分片)**:主分片的完整备份。 - **机制**: - **读写分离**:副本分片可以分担搜索(读)请求,提升查询并发量。 - **故障转移(Failover)**:如果持有主分片的 Node 挂掉,集群会迅速选举一个副本分片升级为新的主分片,保证服务不中断。 ### **3.4 节点角色分化** ![[7-Blog/后端与微服务/assets/Node_角色分化-523fd620.jpg]] 在大型集群中,让每个节点“各司其职”效率更高。 - **Master Node(主节点)**:集群的大脑。负责索引创建/删除、维护集群状态(Cluster State)、管理节点加入/退出。 - **Data Node(数据节点)**:苦力。负责存储数据(Shard),执行耗费资源的 CRUD 和聚合操作。 - **Coordinate Node(协调节点)**:前台/路由。接收客户端请求,分发给 Data Node,并汇总最终结果(Scatter-Gather)。 ### **3.5 去中心化协调机制** ![[7-Blog/后端与微服务/assets/去中心化协调机制-0c561c9f.jpg]] - **Raft 算法衍生**:ES 内部实现了一套基于 Raft 改进的共识算法(Zen Discovery)。 - **作用**: - 保证集群中所有节点对“集群状态”的认知是一致的。 - 实现 Master 节点的选举。 - **故障检测**:节点之间互相 Ping,感知是否有节点掉线。 ## **四、核心流程详解** ### **4.1 写入流程** ES 的写入流程设计权衡了**数据安全性**与**写入高吞吐**。 **1. 路由与转发** - 客户端向任意节点发送写入请求(该节点暂时成为**协调节点**)。 - 协调节点使用**路由算法**确定数据所属的主分片(Primary Shard)位置: - `shard = hash(routing) % number_of_primary_shards` - *注:*`routing` *默认是文档* `_id`*。* - 协调节点将请求转发给持有该主分片的 Data Node。 **2. 主分片写入 (Primary Operation)** - **写入内存缓冲区 (Memory Buffer)**:数据先写入内存 buffer,此时数据**不可被搜索**。 - **写入 Translog (Transaction Log)**:同时追加写入 Translog 文件(顺序写磁盘),防止断电丢失数据。 - **Refresh (准实时关键步骤)**:默认每 1 秒,ES 将 buffer 中的数据生成一个新的 **Segment** 文件(此时建立倒排索引),并清空 buffer。**一旦生成 Segment,数据即可被搜索**。这就是 ES 被称为“准实时(Near Real-Time, NRT)”的原因。 **3. 同步副本 (Replication)** - 主分片写入成功后,并行将请求发送给所有的 **副本分片 (Replica Shards)**。 - 副本分片执行相同的写入逻辑。 **4. 响应客户端** - 当所有在 **ISR (In-Sync Replicas)** 列表中的副本都反馈写入成功后,主分片向协调节点报告成功。 - 协调节点向客户端返回“写入完成”。 > **技术深挖:Flush 操作** Translog 不会无限增长。当 Translog 达到阈值或每隔 30 分钟,ES 会触发 **Flush** 操作: > > 1. 强制执行 Refresh。 > 2. 将所有内存中的 Segment 强制 `fsync` 刷入物理磁盘。 > 3. 清空 Translog。 *这保证了数据的持久化存储。* ![[7-Blog/后端与微服务/assets/ES_写入流程详解-4fcf1ad9.jpg]] ### **4.2 搜索流程(Query Then Fetch)** ![[7-Blog/后端与微服务/assets/ES_搜索流程详解-f639461c.jpg]] 搜索比写入复杂,因为数据分散在多个分片上,必须通过“两阶段”策略来整合结果,以避免网络带宽的巨大浪费。 **阶段一:查询阶段 (Query Phase)** - **请求分发**:客户端向**协调节点**发送搜索请求。协调节点根据请求(是否有 routing 参数)将请求广播到所有相关分片(主分片或副本分片均可,负载均衡)。 - **本地检索**:每个分片在本地 Lucene 中执行搜索: 1. 利用倒排索引筛选匹配文档。 2. 利用 Doc Values 进行排序和打分。 3. **关键点**:分片**仅返回**文档 ID、相关性算分 (\_score) 和排序值给协调节点,**不返回**文档的完整内容 (`_source`)。 - **全局排序**:协调节点收到所有分片返回的轻量级列表(例如每分片前 10 条),在内存中进行**全局归并排序**,选出最终的 Top N 文档 ID。 **阶段二:获取阶段 (Fetch Phase)** - **精确定位**:协调节点知道了最终需要哪几个文档,以及它们位于哪个分片。 - **抓取内容**:协调节点向相关分片发送 `Multi-Get` 请求,只索取这 Top N 文档的完整内容 (`_source` / `Stored Fields`)。 - **返回结果**:分片返回文档详情,协调节点拼装最终 JSON,响应给客户端。 > **性能隐患:深度分页 (Deep Pagination)** 如果查询 `from=10000, size=10`: > > - 每个分片都必须查询出前 10010 条记录。 > - 假设有 5 个分片,协调节点需要接收 `5 * 10010 = 50050` 条记录的 ID,并在内存中排序,最后只取 10 条。 > - **后果**:内存爆炸,CPU 飙升。 > - **对策**:避免深分页,使用 `Search After` 或 `Scroll` API。 ### **4.3 搜索流程总结图** ![[7-Blog/后端与微服务/assets/搜索流程数据结构使用-f707625b.jpg]] ## **五、核心概念对比与 Type 演变** ### **5.1 核心概念对比** ``` ┌─────────────────────────────────────────────────────────────────────────┐ │ ES 与关系型数据库概念对比 │ ├───────────────────┬─────────────────────┬───────────────────────────────┤ │ 关系型数据库 │ Elasticsearch │ 说明 │ ├───────────────────┼─────────────────────┼───────────────────────────────┤ │ Database │ Cluster │ 数据库/集群 │ │ Table │ Index │ 表/索引 │ │ Row │ Document │ 行/文档 │ │ Column │ Field │ 列/字段 │ │ Schema │ Mapping │ 表结构/映射 │ │ Index │ Inverted Index │ 索引/倒排索引 │ │ SQL │ Query DSL │ 查询语言 │ └───────────────────┴─────────────────────┴───────────────────────────────┘ ``` ### **5.2 Type 演变历史** ![[7-Blog/后端与微服务/assets/Type_概念的演变-da19af3e.jpg]] - **5.x 及以前**:允许一个 Index 下存在多个 Type(类比 Table),但本质上底层字段是扁平化混在一起的,导致数据稀疏(Sparse)问题,影响压缩效率和性能。 - **6.x**:强制规定一个 Index 只能有一个 Type,通常默认名为 `doc`。 - **7.x**:Type 概念被彻底废弃(默认为 `_doc`),API 中 URL 的 type 参数变为可选。 - **8.x**:彻底移除 Type 概念。 - **结论**:现在设计索引时,严格遵循 **“一个 Index 对应一类业务数据”** 的原则。 **为什么移除?** - 映射爆炸(Mapping Explosion) - 在使用多类型时,如果不同类型之间有大量不同的字段,这会导致映射的数量急剧增加,进而引发映射爆炸问题。映射爆炸不仅会消耗大量的内存资源,还会降低 Elasticsearch 的性能,尤其是在处理大量数据时 - 字段名冲突 - 在同一个index的不同type中,如果有相同名称但映射类型不同的字段,会造成字段名冲突。这是因为 Elasticsearch 在内部是将这些字段扁平化处理的,而不同类型的相同名称字段可能会导致数据解析和查询时的混乱 - 我们可以和关系型数据库来对比,在同一个数据库中,这些不同的表,可以有名称相同但类型不同的字段。而在 Elasticsearch 同一个index的不同type中,如果有不同document的字段名相同,但是类型不同,就会报错 - 综上所述,Elasticsearch 从 7.x 版本开始废弃类型的主要目的是为了提升系统的性能、避免映射爆炸和字段冲突的问题,以及简化数据模型的设计和管理。这一改变反映了 Elasticsearch 对于提高性能、可维护性和用户体验的持续追求 ## **六、核心功能** ### **6.1 搜索能力矩阵** ![[7-Blog/后端与微服务/assets/ES_搜索能力矩阵-cd6c4c21.jpg]] ### **6.2 聚合分析** ![[7-Blog/后端与微服务/assets/聚合分析类型2-1609a946.jpg]] ## **七、应用场景** ![[7-Blog/后端与微服务/assets/Elasticsearch_应用场景-39904ef5.jpg]] 1. **日志与监控(ELK Stack)**:收集服务器日志、应用 Error 日志,快速定位故障(Logstash/Beats + ES + Kibana)。 2. **站内搜索**:电商商品搜索、论坛帖子搜索、企业知识库检索。 3. **大屏可视化/BI**:实时统计大盘数据(如双11大屏),利用聚合功能快速出报表。 4. **地理位置服务(LBS)**:查询“附近的酒店”、“方圆5公里内的订单”(Geo-point/Geo-shape)。 ## **八、如何使用 Elasticsearch** ### **8.1 Spring Boot 集成** 通常有两种主流方式: 1. **Spring Data Elasticsearch**:封装程度极高,类似 JPA/MyBatis-Plus,通过 Repository 接口操作。 - *优点*:开发极快,代码简洁。 - *缺点*:灵活性稍差,对复杂 DSL 和版本兼容性控制不如原生客户端细致。 2. **RestHighLevelClient (官方推荐/传统)**:基于 HTTP 的原生客户端封装。 - *优点*:完全覆盖官方 API,灵活,可控性强。 - *注意*:ES 7.15+ 后官方推出了新的 `Elasticsearch Java API Client`,但 `RestHighLevelClient` 依然在存量系统中广泛使用。**本文基于此方案进行封装。** ``` org.springframework.boot spring-boot-starter-data-elasticsearch ``` ```yaml # application.yml 配置 spring: elasticsearch: uris: http://localhost:9200 username: elastic password: password ``` 此type的作用就是为了兼容6.x、7.x中的type概念,默认是关闭 ### **8.2 原生 API 操作的复杂性** 直接使用 `RestHighLevelClient` 会面临大量样板代码: - 构建 `SearchSourceBuilder`、`BoolQueryBuilder` 极其繁琐。 - 需要手动处理 `IOException`。 - 响应结果解析(Parse)需要从 JSON 层层剥离,非常痛苦。 - 连接管理和配置分散。 ```java // 原生 ES API 操作示例 - 复杂且繁琐 public SearchResponse searchPrograms(String keyword, Integer categoryId) { // 构建查询条件 BoolQueryBuilder boolQuery = QueryBuilders.boolQuery(); if (StringUtils.isNotBlank(keyword)) { boolQuery.must(QueryBuilders.matchQuery("title", keyword)); } if (categoryId != null) { boolQuery.filter(QueryBuilders.termQuery("categoryId", categoryId)); } // 构建排序 FieldSortBuilder sortBuilder = SortBuilders.fieldSort("showTime") .order(SortOrder.ASC); // 构建搜索源 SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(boolQuery); sourceBuilder.sort(sortBuilder); sourceBuilder.from(0); sourceBuilder.size(10); // 构建搜索请求 SearchRequest searchRequest = new SearchRequest("program-index"); searchRequest.source(sourceBuilder); // 执行搜索 return restHighLevelClient.search(searchRequest, RequestOptions.DEFAULT); } ``` 在平时开发中还是使用关系型数据库更加普遍,对于数据库、表、字段的概念更为熟悉,也更加习惯对表概念的操作。 而在操作Elasticsearch时,提供的api其实是很复杂的,各种操作的对象,如`SearchSourceBuilder FieldSortBuilder` `BoolQueryBuilder` 等等,操作上其实算不上简单,为了解决这个问题,在springboot操作Elasticsearch的基础上,进一步的封装,使用起来贴近于关系型数据库的方式,操作起来更加的容易上手 ## **九、封装设计** ### **9.1 封装目标** - **统一配置**:简化连接参数管理。 - **屏蔽细节**:隐藏繁琐的 Builder 构建过程。 - **简化查询**:通过 Map 或对象传递参数,自动构建 DSL。 - **结果转换**:自动将 ES 的 JSON 结果转为 Java Bean。 - **健壮性**:统一异常处理,防止 ES 波动导致服务崩溃。 ![[7-Blog/后端与微服务/assets/封装设计目标-c8a9c448.jpg]] ### **9.2 封装架构设计** ![[7-Blog/后端与微服务/assets/ES_封装框架架构-9d1b66fd.jpg]] ``` elasticsearch-spring-boot-starter/ ├── src/main/java/com/example/elasticsearch/ │ ├── config/ │ │ ├── ElasticsearchProperties.java # 配置属性 │ │ └── ElasticsearchAutoConfiguration.java # 自动配置 │ ├── core/ │ │ ├── ElasticsearchService.java # 核心服务类 │ │ ├── ElasticsearchIndexService.java # 索引操作服务 │ │ └── ElasticsearchDocumentService.java # 文档操作服务 │ ├── query/ │ │ ├── EsQueryBuilder.java # 查询构建器 │ │ ├── EsSearchRequest.java # 搜索请求封装 │ │ └── EsHighlightConfig.java # 高亮配置 │ ├── result/ │ │ ├── PageResult.java # 分页结果 │ │ ├── SearchResult.java # 搜索结果 │ │ └── AggregationResult.java # 聚合结果 │ ├── exception/ │ │ └── ElasticsearchException.java # 自定义异常 │ └── annotation/ │ ├── EsDocument.java # 文档注解 │ └── EsField.java # 字段注解 └── src/main/resources/ └── META-INF/spring.factories ``` ### **9.3 配置类设计** ```java package com.example.elasticsearch.config; import lombok.Data; import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.validation.annotation.Validated; import javax.validation.constraints.Min; import javax.validation.constraints.NotEmpty; import java.util.List; /** * Elasticsearch 配置属性类 * 支持集群配置、连接池、认证、重试等 */ @Data @Validated @ConfigurationProperties(prefix = "elasticsearch") public class ElasticsearchProperties { /** * 是否启用 ES */ private Boolean enabled = true; /** * ES 节点地址列表(支持集群) */ @NotEmpty(message = "ES 节点地址不能为空") private List nodes = List.of("localhost:9200"); /** * 用户名(可选) */ private String username; /** * 密码(可选) */ private String password; /** * 协议:http 或 https */ private String scheme = "http"; /** * 是否启用 Type(兼容 ES 6.x 版本) */ private Boolean enableType = false; /** * 默认 Type 名称 */ private String defaultType = "_doc"; /** * 连接超时时间(毫秒) */ @Min(value = 1000, message = "连接超时时间不能小于1000ms") private Integer connectTimeout = 5000; /** * Socket 超时时间(毫秒) */ @Min(value = 1000, message = "Socket超时时间不能小于1000ms") private Integer socketTimeout = 30000; /** * 请求超时时间(毫秒) */ private Integer connectionRequestTimeout = 5000; /** * 最大连接数 */ @Min(value = 1, message = "最大连接数不能小于1") private Integer maxConnTotal = 100; /** * 每个路由的最大连接数 */ @Min(value = 1, message = "每个路由最大连接数不能小于1") private Integer maxConnPerRoute = 50; /** * 重试次数 */ @Min(value = 0, message = "重试次数不能为负数") private Integer retryTimes = 3; /** * 重试间隔(毫秒) */ private Long retryInterval = 1000L; /** * 是否开启嗅探器 */ private Boolean enableSniffer = false; /** * 嗅探间隔时间(毫秒) */ private Long snifferInterval = 60000L; /** * 批量操作每批大小 */ @Min(value = 100, message = "批量操作每批大小不能小于100") private Integer bulkBatchSize = 1000; /** * 批量操作刷新策略:immediate, wait_for, none */ private String bulkRefreshPolicy = "none"; /** * 是否打印 DSL 日志 */ private Boolean printDsl = false; /** * 慢查询阈值(毫秒),超过此值记录警告日志 */ private Long slowQueryThreshold = 3000L; } ``` ### **9.4 自动配置类** ```java package com.example.elasticsearch.config; import com.example.elasticsearch.core.ElasticsearchDocumentService; import com.example.elasticsearch.core.ElasticsearchIndexService; import com.example.elasticsearch.core.ElasticsearchService; import lombok.extern.slf4j.Slf4j; import org.apache.http.HttpHost; import org.apache.http.auth.AuthScope; import org.apache.http.auth.UsernamePasswordCredentials; import org.apache.http.client.CredentialsProvider; import org.apache.http.impl.client.BasicCredentialsProvider; import org.elasticsearch.client.RestClient; import org.elasticsearch.client.RestClientBuilder; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.client.sniff.Sniffer; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.util.StringUtils; import javax.annotation.PreDestroy; import java.util.List; import java.util.stream.Collectors; /** * Elasticsearch 自动配置类 */ @Slf4j @Configuration @EnableConfigurationProperties(ElasticsearchProperties.class) @ConditionalOnProperty(prefix = "elasticsearch", name = "enabled", havingValue = "true", matchIfMissing = true) public class ElasticsearchAutoConfiguration { private RestHighLevelClient restHighLevelClient; private Sniffer sniffer; @Bean @ConditionalOnMissingBean public RestHighLevelClient restHighLevelClient(ElasticsearchProperties properties) { // 解析节点地址 List httpHosts = properties.getNodes().stream() .map(node -> { String[] parts = node.split(":"); String host = parts[0]; int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 9200; return new HttpHost(host, port, properties.getScheme()); }) .collect(Collectors.toList()); RestClientBuilder builder = RestClient.builder( httpHosts.toArray(new HttpHost[0]) ); // 设置请求配置 builder.setRequestConfigCallback(requestConfigBuilder -> requestConfigBuilder .setConnectTimeout(properties.getConnectTimeout()) .setSocketTimeout(properties.getSocketTimeout()) .setConnectionRequestTimeout(properties.getConnectionRequestTimeout()) ); // 设置 HTTP 客户端配置 builder.setHttpClientConfigCallback(httpClientBuilder -> { // 设置连接池 httpClientBuilder.setMaxConnTotal(properties.getMaxConnTotal()); httpClientBuilder.setMaxConnPerRoute(properties.getMaxConnPerRoute()); // 设置认证 if (StringUtils.hasText(properties.getUsername()) && StringUtils.hasText(properties.getPassword())) { CredentialsProvider credentialsProvider = new BasicCredentialsProvider(); credentialsProvider.setCredentials( AuthScope.ANY, new UsernamePasswordCredentials( properties.getUsername(), properties.getPassword() ) ); httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider); } return httpClientBuilder; }); // 设置失败重试策略 builder.setFailureListener(new RestClient.FailureListener() { @Override public void onFailure(org.elasticsearch.client.Node node) { log.warn("ES 节点 [{}] 连接失败", node.getHost()); } }); restHighLevelClient = new RestHighLevelClient(builder); // 启用嗅探器 if (properties.getEnableSniffer()) { sniffer = Sniffer.builder(restHighLevelClient.getLowLevelClient()) .setSniffIntervalMillis(properties.getSnifferInterval().intValue()) .build(); log.info("ES 嗅探器已启用,间隔: {}ms", properties.getSnifferInterval()); } log.info("ES 客户端初始化成功, 节点: {}", properties.getNodes()); return restHighLevelClient; } @Bean @ConditionalOnMissingBean public ElasticsearchService elasticsearchService( RestHighLevelClient client, ElasticsearchProperties properties) { return new ElasticsearchService(client, properties); } @Bean @ConditionalOnMissingBean public ElasticsearchIndexService elasticsearchIndexService( RestHighLevelClient client, ElasticsearchProperties properties) { return new ElasticsearchIndexService(client, properties); } @Bean @ConditionalOnMissingBean public ElasticsearchDocumentService elasticsearchDocumentService( RestHighLevelClient client, ElasticsearchProperties properties) { return new ElasticsearchDocumentService(client, properties); } @PreDestroy public void destroy() { try { if (sniffer != null) { sniffer.close(); log.info("ES 嗅探器已关闭"); } if (restHighLevelClient != null) { restHighLevelClient.close(); log.info("ES 客户端已关闭"); } } catch (Exception e) { log.error("关闭 ES 客户端失败", e); } } } ``` ### **9.5 自定义异常类** ```java package com.example.elasticsearch.exception; import lombok.Getter; /** * Elasticsearch 自定义异常 */ @Getter public class ElasticsearchException extends RuntimeException { private static final long serialVersionUID = 1L; /** * 索引名称 */ private String index; /** * 操作类型 */ private String operation; /** * 错误码 */ private String errorCode; /** * 是否可重试 */ private boolean retryable; public ElasticsearchException(String message) { super(message); this.retryable = false; } public ElasticsearchException(String message, Throwable cause) { super(message, cause); this.retryable = isRetryableException(cause); } public ElasticsearchException(String operation, String index, String message) { super(String.format("[%s] 索引 [%s] 操作失败: %s", operation, index, message)); this.operation = operation; this.index = index; this.errorCode = operation + "_ERROR"; } public ElasticsearchException(String operation, String index, Throwable cause) { super(String.format("[%s] 索引 [%s] 操作失败: %s", operation, index, cause.getMessage()), cause); this.operation = operation; this.index = index; this.errorCode = operation + "_ERROR"; this.retryable = isRetryableException(cause); } /** * 判断异常是否可重试 */ private boolean isRetryableException(Throwable cause) { if (cause == null) { return false; } String message = cause.getMessage(); if (message == null) { return false; } // 可重试的异常类型 return message.contains("Connection refused") || message.contains("Connection reset") || message.contains("Connection timed out") || message.contains("Read timed out") || message.contains("No route to host") || message.contains("Service Unavailable") || message.contains("circuit_breaking_exception"); } /** * 创建索引不存在异常 */ public static ElasticsearchException indexNotFound(String index) { ElasticsearchException ex = new ElasticsearchException( "INDEX_NOT_FOUND", index, "索引不存在"); ex.errorCode = "INDEX_NOT_FOUND"; return ex; } /** * 创建文档不存在异常 */ public static ElasticsearchException documentNotFound(String index, String id) { ElasticsearchException ex = new ElasticsearchException( "DOCUMENT_NOT_FOUND", index, String.format("文档 [%s] 不存在", id)); ex.errorCode = "DOCUMENT_NOT_FOUND"; return ex; } /** * 创建参数校验异常 */ public static ElasticsearchException invalidParameter(String message) { ElasticsearchException ex = new ElasticsearchException(message); ex.errorCode = "INVALID_PARAMETER"; return ex; } } ``` ### **9.6 结果封装类** ```java package com.example.elasticsearch.result; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import java.util.Collections; import java.util.List; /** * 分页结果封装 */ @Data @NoArgsConstructor @AllArgsConstructor public class PageResult { /** * 数据列表 */ private List records; /** * 总记录数 */ private Long total; /** * 当前页码 */ private Integer pageNum; /** * 每页大小 */ private Integer pageSize; /** * 总页数 */ private Integer totalPages; /** * 是否有下一页 */ private Boolean hasNext; /** * 是否有上一页 */ private Boolean hasPrevious; public PageResult(List records, Long total, Integer pageNum, Integer pageSize) { this.records = records; this.total = total; this.pageNum = pageNum; this.pageSize = pageSize; this.totalPages = (int) Math.ceil((double) total / pageSize); this.hasNext = pageNum < totalPages; this.hasPrevious = pageNum > 1; } /** * 空结果 */ public static PageResult empty(Integer pageNum, Integer pageSize) { return new PageResult<>(Collections.emptyList(), 0L, pageNum, pageSize); } } ``` ```java package com.example.elasticsearch.result; import lombok.Data; import java.util.List; import java.util.Map; /** * 搜索结果封装(包含高亮、评分等信息) */ @Data public class SearchResult { /** * 文档 ID */ private String documentId; /** * 数据对象 */ private T source; /** * 评分 */ private Float score; /** * 高亮字段 */ private Map> highlight; /** * 排序值(用于深度分页) */ private Object[] sortValues; public SearchResult(String documentId, T source) { this.documentId = documentId; this.source = source; } public SearchResult(String documentId, T source, Float score, Map> highlight) { this.documentId = documentId; this.source = source; this.score = score; this.highlight = highlight; } } ``` ```java package com.example.elasticsearch.result; import lombok.Data; import java.util.List; import java.util.Map; /** * 聚合结果封装 */ @Data public class AggregationResult { /** * 聚合名称 */ private String name; /** * 桶数据(terms 聚合) */ private List buckets; /** * 数值(sum、avg、max、min 等) */ private Double value; @Data public static class BucketData { private String key; private Long docCount; private Map subAggregations; } } ``` ### **9.7 查询构建器** ```java package com.example.elasticsearch.query; import lombok.Data; import org.elasticsearch.index.query.*; import org.elasticsearch.search.aggregations.AggregationBuilder; import org.elasticsearch.search.aggregations.AggregationBuilders; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder; import org.elasticsearch.search.sort.SortBuilder; import org.elasticsearch.search.sort.SortBuilders; import org.elasticsearch.search.sort.SortOrder; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import java.util.*; /** * ES 查询构建器 * 使用链式调用构建复杂查询 */ @Data public class EsQueryBuilder { /** * 索引名称 */ private String indexName; /** * Bool 查询条件 */ private BoolQueryBuilder boolQuery; /** * 排序条件 */ private List> sorts; /** * 高亮配置 */ private HighlightBuilder highlightBuilder; /** * 聚合配置 */ private List aggregations; /** * 分页参数 */ private Integer from; private Integer size; /** * 返回字段(包含) */ private String[] includes; /** * 排除字段 */ private String[] excludes; /** * Search After(深度分页) */ private Object[] searchAfter; /** * 是否追踪总数 */ private Boolean trackTotalHits = true; private EsQueryBuilder() { this.boolQuery = QueryBuilders.boolQuery(); this.sorts = new ArrayList<>(); this.aggregations = new ArrayList<>(); } /** * 创建构建器 */ public static EsQueryBuilder builder(String indexName) { EsQueryBuilder builder = new EsQueryBuilder(); builder.indexName = indexName; return builder; } // ==================== Must 条件(必须匹配)==================== /** * 精确匹配(term) */ public EsQueryBuilder term(String field, Object value) { if (value != null) { boolQuery.must(QueryBuilders.termQuery(field, value)); } return this; } /** * 多值匹配(terms) */ public EsQueryBuilder terms(String field, Collection values) { if (!CollectionUtils.isEmpty(values)) { boolQuery.must(QueryBuilders.termsQuery(field, values)); } return this; } /** * 全文匹配(match) */ public EsQueryBuilder match(String field, Object value) { if (value != null && StringUtils.hasText(value.toString())) { boolQuery.must(QueryBuilders.matchQuery(field, value)); } return this; } /** * 短语匹配(match_phrase) */ public EsQueryBuilder matchPhrase(String field, Object value) { if (value != null && StringUtils.hasText(value.toString())) { boolQuery.must(QueryBuilders.matchPhraseQuery(field, value)); } return this; } /** * 多字段匹配(multi_match) */ public EsQueryBuilder multiMatch(Object value, String... fields) { if (value != null && StringUtils.hasText(value.toString()) && fields.length > 0) { boolQuery.must(QueryBuilders.multiMatchQuery(value, fields)); } return this; } /** * 前缀匹配(prefix) */ public EsQueryBuilder prefix(String field, String prefix) { if (StringUtils.hasText(prefix)) { boolQuery.must(QueryBuilders.prefixQuery(field, prefix)); } return this; } /** * 通配符匹配(wildcard) */ public EsQueryBuilder wildcard(String field, String pattern) { if (StringUtils.hasText(pattern)) { boolQuery.must(QueryBuilders.wildcardQuery(field, pattern)); } return this; } /** * 范围查询 */ public EsQueryBuilder range(String field, Object gte, Object lte) { RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field); if (gte != null) { rangeQuery.gte(gte); } if (lte != null) { rangeQuery.lte(lte); } if (gte != null || lte != null) { boolQuery.must(rangeQuery); } return this; } /** * 范围查询(大于) */ public EsQueryBuilder gt(String field, Object value) { if (value != null) { boolQuery.must(QueryBuilders.rangeQuery(field).gt(value)); } return this; } /** * 范围查询(大于等于) */ public EsQueryBuilder gte(String field, Object value) { if (value != null) { boolQuery.must(QueryBuilders.rangeQuery(field).gte(value)); } return this; } /** * 范围查询(小于) */ public EsQueryBuilder lt(String field, Object value) { if (value != null) { boolQuery.must(QueryBuilders.rangeQuery(field).lt(value)); } return this; } /** * 范围查询(小于等于) */ public EsQueryBuilder lte(String field, Object value) { if (value != null) { boolQuery.must(QueryBuilders.rangeQuery(field).lte(value)); } return this; } /** * 存在字段查询 */ public EsQueryBuilder exists(String field) { boolQuery.must(QueryBuilders.existsQuery(field)); return this; } // ==================== Filter 条件(过滤,不计算评分)==================== /** * Filter - 精确匹配 */ public EsQueryBuilder filterTerm(String field, Object value) { if (value != null) { boolQuery.filter(QueryBuilders.termQuery(field, value)); } return this; } /** * Filter - 多值匹配 */ public EsQueryBuilder filterTerms(String field, Collection values) { if (!CollectionUtils.isEmpty(values)) { boolQuery.filter(QueryBuilders.termsQuery(field, values)); } return this; } /** * Filter - 范围查询 */ public EsQueryBuilder filterRange(String field, Object gte, Object lte) { RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field); if (gte != null) { rangeQuery.gte(gte); } if (lte != null) { rangeQuery.lte(lte); } if (gte != null || lte != null) { boolQuery.filter(rangeQuery); } return this; } // ==================== Should 条件(或条件)==================== /** * Should - 至少匹配一个 */ public EsQueryBuilder should(QueryBuilder... queries) { for (QueryBuilder query : queries) { boolQuery.should(query); } return this; } /** * Should - 多字段或查询 */ public EsQueryBuilder shouldMatch(String value, String... fields) { if (StringUtils.hasText(value) && fields.length > 0) { for (String field : fields) { boolQuery.should(QueryBuilders.matchQuery(field, value)); } boolQuery.minimumShouldMatch(1); } return this; } /** * 设置最小 should 匹配数 */ public EsQueryBuilder minimumShouldMatch(int count) { boolQuery.minimumShouldMatch(count); return this; } // ==================== MustNot 条件(必须不匹配)==================== /** * MustNot - 精确匹配 */ public EsQueryBuilder mustNotTerm(String field, Object value) { if (value != null) { boolQuery.mustNot(QueryBuilders.termQuery(field, value)); } return this; } /** * MustNot - 多值匹配 */ public EsQueryBuilder mustNotTerms(String field, Collection values) { if (!CollectionUtils.isEmpty(values)) { boolQuery.mustNot(QueryBuilders.termsQuery(field, values)); } return this; } // ==================== 嵌套查询 ==================== /** * 嵌套查询 */ public EsQueryBuilder nested(String path, QueryBuilder query) { boolQuery.must(QueryBuilders.nestedQuery(path, query, org.apache.lucene.search.join.ScoreMode.Avg)); return this; } /** * 添加自定义查询条件 */ public EsQueryBuilder must(QueryBuilder query) { boolQuery.must(query); return this; } /** * 添加自定义过滤条件 */ public EsQueryBuilder filter(QueryBuilder query) { boolQuery.filter(query); return this; } // ==================== 排序 ==================== /** * 添加排序 */ public EsQueryBuilder sort(String field, SortOrder order) { sorts.add(SortBuilders.fieldSort(field).order(order)); return this; } /** * 按评分排序 */ public EsQueryBuilder sortByScore(SortOrder order) { sorts.add(SortBuilders.scoreSort().order(order)); return this; } /** * 多字段排序 */ public EsQueryBuilder sorts(Map sortMap) { sortMap.forEach((field, order) -> sorts.add(SortBuilders.fieldSort(field).order(order))); return this; } // ==================== 分页 ==================== /** * 设置分页 */ public EsQueryBuilder page(int pageNum, int pageSize) { this.from = (pageNum - 1) * pageSize; this.size = pageSize; return this; } /** * 设置起始位置和大小 */ public EsQueryBuilder fromSize(int from, int size) { this.from = from; this.size = size; return this; } /** * 深度分页(Search After) */ public EsQueryBuilder searchAfter(Object[] values) { this.searchAfter = values; return this; } // ==================== 高亮 ==================== /** * 添加高亮字段 */ public EsQueryBuilder highlight(String... fields) { if (fields.length > 0) { highlightBuilder = new HighlightBuilder(); for (String field : fields) { highlightBuilder.field(field); } highlightBuilder.preTags(""); highlightBuilder.postTags(""); } return this; } /** * 自定义高亮配置 */ public EsQueryBuilder highlight(String preTag, String postTag, String... fields) { if (fields.length > 0) { highlightBuilder = new HighlightBuilder(); for (String field : fields) { highlightBuilder.field(field); } highlightBuilder.preTags(preTag); highlightBuilder.postTags(postTag); } return this; } // ==================== 聚合 ==================== /** * Terms 聚合 */ public EsQueryBuilder termsAggregation(String name, String field, int size) { aggregations.add(AggregationBuilders.terms(name).field(field).size(size)); return this; } /** * Sum 聚合 */ public EsQueryBuilder sumAggregation(String name, String field) { aggregations.add(AggregationBuilders.sum(name).field(field)); return this; } /** * Avg 聚合 */ public EsQueryBuilder avgAggregation(String name, String field) { aggregations.add(AggregationBuilders.avg(name).field(field)); return this; } /** * Max 聚合 */ public EsQueryBuilder maxAggregation(String name, String field) { aggregations.add(AggregationBuilders.max(name).field(field)); return this; } /** * Min 聚合 */ public EsQueryBuilder minAggregation(String name, String field) { aggregations.add(AggregationBuilders.min(name).field(field)); return this; } /** * 日期直方图聚合 */ public EsQueryBuilder dateHistogramAggregation(String name, String field, String interval) { aggregations.add(AggregationBuilders.dateHistogram(name) .field(field) .calendarInterval(new org.elasticsearch.search.aggregations.bucket .histogram.DateHistogramInterval(interval))); return this; } /** * 添加自定义聚合 */ public EsQueryBuilder aggregation(AggregationBuilder aggregation) { aggregations.add(aggregation); return this; } // ==================== 返回字段 ==================== /** * 指定返回字段 */ public EsQueryBuilder includes(String... fields) { this.includes = fields; return this; } /** * 排除字段 */ public EsQueryBuilder excludes(String... fields) { this.excludes = fields; return this; } /** * 是否追踪总数 */ public EsQueryBuilder trackTotalHits(boolean track) { this.trackTotalHits = track; return this; } // ==================== 构建 ==================== /** * 构建 SearchSourceBuilder */ public SearchSourceBuilder build() { SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); // 查询条件 sourceBuilder.query(boolQuery); // 分页 if (from != null) { sourceBuilder.from(from); } if (size != null) { sourceBuilder.size(size); } // 排序 for (SortBuilder sort : sorts) { sourceBuilder.sort(sort); } // 高亮 if (highlightBuilder != null) { sourceBuilder.highlighter(highlightBuilder); } // 聚合 for (AggregationBuilder aggregation : aggregations) { sourceBuilder.aggregation(aggregation); } // 返回字段 if (includes != null || excludes != null) { sourceBuilder.fetchSource(includes, excludes); } // Search After if (searchAfter != null) { sourceBuilder.searchAfter(searchAfter); } // 追踪总数 sourceBuilder.trackTotalHits(trackTotalHits); return sourceBuilder; } } ``` ### **9.8 重试工具类** ```java package com.example.elasticsearch.util; import com.example.elasticsearch.exception.ElasticsearchException; import lombok.extern.slf4j.Slf4j; import java.util.concurrent.Callable; import java.util.function.Predicate; /** * 重试工具类 */ @Slf4j public class RetryUtil { /** * 执行带重试的操作 * * @param callable 要执行的操作 * @param maxRetries 最大重试次数 * @param retryInterval 重试间隔(毫秒) * @param retryOn 判断是否需要重试的条件 * @param operationName 操作名称(用于日志) * @return 操作结果 */ public static T executeWithRetry( Callable callable, int maxRetries, long retryInterval, Predicate retryOn, String operationName) { Exception lastException = null; for (int attempt = 0; attempt <= maxRetries; attempt++) { try { return callable.call(); } catch (Exception e) { lastException = e; // 判断是否需要重试 if (attempt < maxRetries && retryOn.test(e)) { log.warn("[{}] 操作失败,第 {}/{} 次重试,错误: {}", operationName, attempt + 1, maxRetries, e.getMessage()); try { Thread.sleep(retryInterval * (attempt + 1)); // 指数退避 } catch (InterruptedException ie) { Thread.currentThread().interrupt(); throw new ElasticsearchException("重试被中断", ie); } } else { break; } } } log.error("[{}] 操作失败,已达最大重试次数", operationName, lastException); if (lastException instanceof ElasticsearchException) { throw (ElasticsearchException) lastException; } throw new ElasticsearchException(operationName + " 操作失败", lastException); } /** * 执行带重试的操作(无返回值) */ public static void executeWithRetry( Runnable runnable, int maxRetries, long retryInterval, Predicate retryOn, String operationName) { executeWithRetry(() -> { runnable.run(); return null; }, maxRetries, retryInterval, retryOn, operationName); } /** * 默认的重试条件判断 */ public static Predicate defaultRetryCondition() { return e -> { if (e instanceof ElasticsearchException) { return ((ElasticsearchException) e).isRetryable(); } String message = e.getMessage(); if (message == null) { return false; } return message.contains("Connection") || message.contains("timed out") || message.contains("Unavailable"); }; } } ``` ### **9.9 参数校验工具类** ```java package com.example.elasticsearch.util; import com.example.elasticsearch.exception.ElasticsearchException; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import java.util.Collection; /** * 参数校验工具类 */ public class ParamValidator { private ParamValidator() {} /** * 校验索引名称 */ public static void validateIndexName(String indexName) { if (!StringUtils.hasText(indexName)) { throw ElasticsearchException.invalidParameter("索引名称不能为空"); } if (indexName.contains(" ")) { throw ElasticsearchException.invalidParameter("索引名称不能包含空格"); } if (!indexName.equals(indexName.toLowerCase())) { throw ElasticsearchException.invalidParameter("索引名称必须小写"); } } /** * 校验文档ID */ public static void validateDocumentId(String documentId) { if (!StringUtils.hasText(documentId)) { throw ElasticsearchException.invalidParameter("文档ID不能为空"); } } /** * 校验分页参数 */ public static void validatePageParam(int pageNum, int pageSize) { if (pageNum < 1) { throw ElasticsearchException.invalidParameter("页码必须大于0"); } if (pageSize < 1 || pageSize > 10000) { throw ElasticsearchException.invalidParameter("每页大小必须在1-10000之间"); } // ES 默认限制 from + size <= 10000 if ((long) (pageNum - 1) * pageSize + pageSize > 10000) { throw ElasticsearchException.invalidParameter( "分页深度超出限制,请使用 searchAfter 方式"); } } /** * 校验批量数据 */ public static void validateBatchData(Collection dataList) { if (CollectionUtils.isEmpty(dataList)) { throw ElasticsearchException.invalidParameter("批量数据不能为空"); } } /** * 校验非空 */ public static void notNull(Object object, String message) { if (object == null) { throw ElasticsearchException.invalidParameter(message); } } /** * 校验字符串非空 */ public static void notBlank(String str, String message) { if (!StringUtils.hasText(str)) { throw ElasticsearchException.invalidParameter(message); } } } ``` ### **9.10 核心 Service 封装** ```java package com.example.elasticsearch.core; import com.alibaba.fastjson.JSON; import com.example.elasticsearch.config.ElasticsearchProperties; import com.example.elasticsearch.exception.ElasticsearchException; import com.example.elasticsearch.query.EsQueryBuilder; import com.example.elasticsearch.result.AggregationResult; import com.example.elasticsearch.result.PageResult; import com.example.elasticsearch.result.ScrollResult; import com.example.elasticsearch.result.SearchResult; import com.example.elasticsearch.util.ParamValidator; import com.example.elasticsearch.util.RetryUtil; import lombok.extern.slf4j.Slf4j; import org.elasticsearch.action.search.*; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.common.text.Text; import org.elasticsearch.common.unit.TimeValue; import org.elasticsearch.index.query.QueryBuilder; import org.elasticsearch.search.SearchHit; import org.elasticsearch.search.SearchHits; import org.elasticsearch.search.aggregations.Aggregation; import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.metrics.*; import org.elasticsearch.search.builder.SearchSourceBuilder; import org.elasticsearch.search.fetch.subphase.highlight.HighlightField; import java.io.IOException; import java.util.*; import java.util.stream.Collectors; /** * Elasticsearch 核心搜索服务 * * 特性: * - 完善的参数校验 * - 自动重试机制 * - 慢查询日志 * - 滚动查询支持 * - 空值安全处理 */ @Slf4j public class ElasticsearchService { private final RestHighLevelClient client; private final ElasticsearchProperties properties; public ElasticsearchService(RestHighLevelClient client, ElasticsearchProperties properties) { this.client = client; this.properties = properties; } // ==================== 基础查询 ==================== /** * 简单查询 - 根据单个字段精确匹配 * * @param indexName 索引名称 * @param field 字段名 * @param value 字段值 * @param clazz 返回类型 * @return 匹配的文档列表 */ public List query(String indexName, String field, Object value, Class clazz) { ParamValidator.validateIndexName(indexName); ParamValidator.notBlank(field, "查询字段不能为空"); EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName) .term(field, value); return search(queryBuilder, clazz); } /** * 多条件查询 - 根据多个字段匹配 * * @param indexName 索引名称 * @param params 查询参数 (字段名 -> 字段值) * @param clazz 返回类型 * @return 匹配的文档列表 */ public List query(String indexName, Map params, Class clazz) { ParamValidator.validateIndexName(indexName); EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName); if (params != null && !params.isEmpty()) { params.forEach((field, value) -> { if (value != null) { if (value instanceof String && ((String) value).length() > 0) { queryBuilder.match(field, value); } else if (!(value instanceof String)) { queryBuilder.filterTerm(field, value); } } }); } return search(queryBuilder, clazz); } /** * 分页查询 * * @param indexName 索引名称 * @param params 查询参数 * @param pageNum 页码(从1开始) * @param pageSize 每页大小 * @param clazz 返回类型 * @return 分页结果 */ public PageResult queryPage(String indexName, Map params, int pageNum, int pageSize, Class clazz) { ParamValidator.validateIndexName(indexName); ParamValidator.validatePageParam(pageNum, pageSize); EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName) .page(pageNum, pageSize); if (params != null && !params.isEmpty()) { params.forEach((field, value) -> { if (value != null) { if (value instanceof String && ((String) value).length() > 0) { queryBuilder.match(field, value); } else if (!(value instanceof String)) { queryBuilder.filterTerm(field, value); } } }); } return searchPage(queryBuilder, pageNum, pageSize, clazz); } // ==================== 高级查询 ==================== /** * 使用查询构建器执行查询 * * @param queryBuilder 查询构建器 * @param clazz 返回类型 * @return 匹配的文档列表 */ public List search(EsQueryBuilder queryBuilder, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); ParamValidator.notNull(clazz, "返回类型不能为空"); return RetryUtil.executeWithRetry( () -> doSearch(queryBuilder, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH" ); } private List doSearch(EsQueryBuilder queryBuilder, Class clazz) throws IOException { long startTime = System.currentTimeMillis(); try { SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); // 打印 DSL if (properties.getPrintDsl()) { log.info("ES 查询 DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); // 检查响应状态 checkResponseStatus(response, queryBuilder.getIndexName()); return parseHits(response.getHits(), clazz); } finally { logSlowQuery(startTime, "SEARCH", queryBuilder.getIndexName()); } } /** * 分页查询 */ public PageResult searchPage(EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); ParamValidator.validatePageParam(pageNum, pageSize); ParamValidator.notNull(clazz, "返回类型不能为空"); return RetryUtil.executeWithRetry( () -> doSearchPage(queryBuilder, pageNum, pageSize, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH_PAGE" ); } private PageResult doSearchPage(EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class clazz) throws IOException { long startTime = System.currentTimeMillis(); try { queryBuilder.page(pageNum, pageSize); SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); if (properties.getPrintDsl()) { log.info("ES 分页查询 DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); List records = parseHits(response.getHits(), clazz); long total = getTotalHits(response.getHits()); return new PageResult<>(records, total, pageNum, pageSize); } finally { logSlowQuery(startTime, "SEARCH_PAGE", queryBuilder.getIndexName()); } } /** * 带高亮的查询 */ public List> searchWithHighlight(EsQueryBuilder queryBuilder, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); return RetryUtil.executeWithRetry( () -> doSearchWithHighlight(queryBuilder, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH_HIGHLIGHT" ); } private List> doSearchWithHighlight( EsQueryBuilder queryBuilder, Class clazz) throws IOException { long startTime = System.currentTimeMillis(); try { SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); if (properties.getPrintDsl()) { log.info("ES 高亮查询 DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); return parseHitsWithHighlight(response.getHits(), clazz); } finally { logSlowQuery(startTime, "SEARCH_HIGHLIGHT", queryBuilder.getIndexName()); } } /** * 带高亮的分页查询 */ public PageResult> searchPageWithHighlight( EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); ParamValidator.validatePageParam(pageNum, pageSize); return RetryUtil.executeWithRetry( () -> doSearchPageWithHighlight(queryBuilder, pageNum, pageSize, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH_PAGE_HIGHLIGHT" ); } private PageResult> doSearchPageWithHighlight( EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class clazz) throws IOException { long startTime = System.currentTimeMillis(); try { queryBuilder.page(pageNum, pageSize); SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); List> records = parseHitsWithHighlight(response.getHits(), clazz); long total = getTotalHits(response.getHits()); return new PageResult<>(records, total, pageNum, pageSize); } finally { logSlowQuery(startTime, "SEARCH_PAGE_HIGHLIGHT", queryBuilder.getIndexName()); } } // ==================== 滚动查询(大数据量) ==================== /** * 滚动查询 - 初始化 * 适用于导出大量数据的场景 * * @param queryBuilder 查询构建器 * @param scrollTime 滚动上下文保持时间(分钟) * @param size 每批大小 * @param clazz 返回类型 * @return 滚动结果(包含scrollId和首批数据) */ public ScrollResult scrollSearch(EsQueryBuilder queryBuilder, int scrollTime, int size, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); try { queryBuilder.fromSize(0, size); SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); searchRequest.scroll(TimeValue.timeValueMinutes(scrollTime)); if (properties.getPrintDsl()) { log.info("ES 滚动查询初始化 DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); List records = parseHits(response.getHits(), clazz); long total = getTotalHits(response.getHits()); String scrollId = response.getScrollId(); return new ScrollResult<>(scrollId, records, total, records.size() < size); } catch (IOException e) { throw new ElasticsearchException("SCROLL_SEARCH", queryBuilder.getIndexName(), e); } } /** * 滚动查询 - 继续获取下一批 * * @param scrollId 滚动ID * @param scrollTime 滚动上下文保持时间(分钟) * @param size 每批大小 * @param clazz 返回类型 * @return 滚动结果 */ public ScrollResult scrollNext(String scrollId, int scrollTime, int size, Class clazz) { ParamValidator.notBlank(scrollId, "scrollId不能为空"); try { SearchScrollRequest scrollRequest = new SearchScrollRequest(scrollId); scrollRequest.scroll(TimeValue.timeValueMinutes(scrollTime)); SearchResponse response = client.scroll(scrollRequest, RequestOptions.DEFAULT); List records = parseHits(response.getHits(), clazz); String newScrollId = response.getScrollId(); return new ScrollResult<>(newScrollId, records, getTotalHits(response.getHits()), records.size() < size); } catch (IOException e) { throw new ElasticsearchException("滚动查询失败", e); } } /** * 清除滚动上下文 * * @param scrollIds 滚动ID列表 */ public void clearScroll(String... scrollIds) { if (scrollIds == null || scrollIds.length == 0) { return; } try { ClearScrollRequest clearScrollRequest = new ClearScrollRequest(); clearScrollRequest.scrollIds(Arrays.asList(scrollIds)); ClearScrollResponse response = client.clearScroll( clearScrollRequest, RequestOptions.DEFAULT); if (!response.isSucceeded()) { log.warn("清除滚动上下文失败"); } } catch (IOException e) { log.warn("清除滚动上下文异常", e); } } // ==================== 聚合查询 ==================== /** * 聚合查询 */ public Map searchAggregation(EsQueryBuilder queryBuilder) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); return RetryUtil.executeWithRetry( () -> doSearchAggregation(queryBuilder), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH_AGGREGATION" ); } private Map doSearchAggregation( EsQueryBuilder queryBuilder) throws IOException { long startTime = System.currentTimeMillis(); try { queryBuilder.fromSize(0, 0); SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); if (properties.getPrintDsl()) { log.info("ES 聚合查询 DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); return parseAggregations(response.getAggregations()); } finally { logSlowQuery(startTime, "SEARCH_AGGREGATION", queryBuilder.getIndexName()); } } // ==================== Search After(深度分页) ==================== /** * Search After 查询 * 适用于深度分页场景,避免 from+size 的 10000 限制 * * @param queryBuilder 查询构建器 * @param searchAfterValues 上一页最后一条记录的排序值 * @param size 每页大小 * @param clazz 返回类型 * @return 搜索结果列表 */ public List> searchAfter(EsQueryBuilder queryBuilder, Object[] searchAfterValues, int size, Class clazz) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); return RetryUtil.executeWithRetry( () -> doSearchAfter(queryBuilder, searchAfterValues, size, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "SEARCH_AFTER" ); } private List> doSearchAfter( EsQueryBuilder queryBuilder, Object[] searchAfterValues, int size, Class clazz) throws IOException { long startTime = System.currentTimeMillis(); try { queryBuilder.fromSize(0, size); if (searchAfterValues != null && searchAfterValues.length > 0) { queryBuilder.searchAfter(searchAfterValues); } SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); if (properties.getPrintDsl()) { log.info("ES Search After DSL: {}", sourceBuilder.toString()); } SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); return parseHitsWithHighlight(response.getHits(), clazz); } finally { logSlowQuery(startTime, "SEARCH_AFTER", queryBuilder.getIndexName()); } } // ==================== 统计查询 ==================== /** * 统计符合条件的文档数量 */ public long count(EsQueryBuilder queryBuilder) { ParamValidator.validateIndexName(queryBuilder.getIndexName()); return RetryUtil.executeWithRetry( () -> doCount(queryBuilder), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "COUNT" ); } private long doCount(EsQueryBuilder queryBuilder) throws IOException { queryBuilder.fromSize(0, 0); SearchSourceBuilder sourceBuilder = queryBuilder.build(); SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName()); searchRequest.source(sourceBuilder); SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, queryBuilder.getIndexName()); return getTotalHits(response.getHits()); } /** * 判断是否存在符合条件的文档 */ public boolean exists(EsQueryBuilder queryBuilder) { return count(queryBuilder) > 0; } // ==================== 原生查询 ==================== /** * 执行原生 QueryBuilder 查询 * 适用于封装方法无法满足的复杂场景 */ public List searchByQueryBuilder(String indexName, QueryBuilder queryBuilder, int size, Class clazz) { ParamValidator.validateIndexName(indexName); ParamValidator.notNull(queryBuilder, "查询条件不能为空"); try { SearchSourceBuilder sourceBuilder = new SearchSourceBuilder(); sourceBuilder.query(queryBuilder); sourceBuilder.size(size); sourceBuilder.trackTotalHits(true); SearchRequest searchRequest = new SearchRequest(indexName); searchRequest.source(sourceBuilder); SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT); checkResponseStatus(response, indexName); return parseHits(response.getHits(), clazz); } catch (IOException e) { throw new ElasticsearchException("SEARCH_BY_QUERY_BUILDER", indexName, e); } } // ==================== 内部工具方法 ==================== /** * 解析命中结果 */ private List parseHits(SearchHits hits, Class clazz) { if (hits == null || hits.getHits() == null) { return Collections.emptyList(); } List result = new ArrayList<>(); for (SearchHit hit : hits.getHits()) { try { String source = hit.getSourceAsString(); if (source != null && !source.isEmpty()) { T obj = JSON.parseObject(source, clazz); result.add(obj); } } catch (Exception e) { log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage()); } } return result; } /** * 解析带高亮的命中结果 */ private List> parseHitsWithHighlight(SearchHits hits, Class clazz) { if (hits == null || hits.getHits() == null) { return Collections.emptyList(); } List> result = new ArrayList<>(); for (SearchHit hit : hits.getHits()) { try { String source = hit.getSourceAsString(); if (source == null || source.isEmpty()) { continue; } T obj = JSON.parseObject(source, clazz); // 解析高亮 Map> highlightMap = new HashMap<>(); Map highlightFields = hit.getHighlightFields(); if (highlightFields != null && !highlightFields.isEmpty()) { highlightFields.forEach((field, highlightField) -> { if (highlightField.getFragments() != null) { List fragments = Arrays.stream(highlightField.getFragments()) .map(Text::string) .collect(Collectors.toList()); highlightMap.put(field, fragments); } }); } SearchResult searchResult = new SearchResult<>( hit.getId(), obj, hit.getScore(), highlightMap); searchResult.setSortValues(hit.getSortValues()); result.add(searchResult); } catch (Exception e) { log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage()); } } return result; } /** * 解析聚合结果 */ private Map parseAggregations( org.elasticsearch.search.aggregations.Aggregations aggregations) { Map result = new HashMap<>(); if (aggregations == null) { return result; } for (Aggregation aggregation : aggregations) { try { AggregationResult aggResult = new AggregationResult(); aggResult.setName(aggregation.getName()); if (aggregation instanceof Terms) { Terms terms = (Terms) aggregation; List buckets = new ArrayList<>(); for (Terms.Bucket bucket : terms.getBuckets()) { AggregationResult.BucketData bucketData = new AggregationResult.BucketData(); bucketData.setKey(bucket.getKeyAsString()); bucketData.setDocCount(bucket.getDocCount()); buckets.add(bucketData); } aggResult.setBuckets(buckets); } else if (aggregation instanceof Sum) { aggResult.setValue(((Sum) aggregation).getValue()); } else if (aggregation instanceof Avg) { aggResult.setValue(((Avg) aggregation).getValue()); } else if (aggregation instanceof Max) { aggResult.setValue(((Max) aggregation).getValue()); } else if (aggregation instanceof Min) { aggResult.setValue(((Min) aggregation).getValue()); } else if (aggregation instanceof ValueCount) { aggResult.setValue((double) ((ValueCount) aggregation).getValue()); } else if (aggregation instanceof Cardinality) { aggResult.setValue((double) ((Cardinality) aggregation).getValue()); } result.put(aggregation.getName(), aggResult); } catch (Exception e) { log.warn("解析聚合结果失败, name: {}, error: {}", aggregation.getName(), e.getMessage()); } } return result; } /** * 获取总命中数(兼容不同版本) */ private long getTotalHits(SearchHits hits) { if (hits == null || hits.getTotalHits() == null) { return 0L; } return hits.getTotalHits().value; } /** * 检查响应状态 */ private void checkResponseStatus(SearchResponse response, String indexName) { if (response == null) { throw new ElasticsearchException("SEARCH", indexName, "响应为空"); } if (response.isTimedOut()) { log.warn("ES 查询超时, index: {}", indexName); } if (response.getShardFailures() != null && response.getShardFailures().length > 0) { log.warn("ES 查询部分分片失败, index: {}, failures: {}", indexName, response.getShardFailures().length); } } /** * 记录慢查询日志 */ private void logSlowQuery(long startTime, String operation, String indexName) { long elapsed = System.currentTimeMillis() - startTime; if (elapsed >= properties.getSlowQueryThreshold()) { log.warn("[慢查询] 操作: {}, 索引: {}, 耗时: {}ms", operation, indexName, elapsed); } else { log.debug("ES 操作完成, 操作: {}, 索引: {}, 耗时: {}ms", operation, indexName, elapsed); } } } ``` ### **9.11 文档操作服务** ```java package com.example.elasticsearch.core; import com.alibaba.fastjson.JSON; import com.example.elasticsearch.config.ElasticsearchProperties; import com.example.elasticsearch.exception.ElasticsearchException; import com.example.elasticsearch.util.ParamValidator; import com.example.elasticsearch.util.RetryUtil; import lombok.extern.slf4j.Slf4j; import org.elasticsearch.action.DocWriteResponse; import org.elasticsearch.action.bulk.BulkItemResponse; import org.elasticsearch.action.bulk.BulkRequest; import org.elasticsearch.action.bulk.BulkResponse; import org.elasticsearch.action.delete.DeleteRequest; import org.elasticsearch.action.delete.DeleteResponse; import org.elasticsearch.action.get.*; import org.elasticsearch.action.index.IndexRequest; import org.elasticsearch.action.index.IndexResponse; import org.elasticsearch.action.support.WriteRequest; import org.elasticsearch.action.update.UpdateRequest; import org.elasticsearch.action.update.UpdateResponse; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.common.xcontent.XContentType; import org.elasticsearch.index.query.QueryBuilder; import org.elasticsearch.index.reindex.BulkByScrollResponse; import org.elasticsearch.index.reindex.DeleteByQueryRequest; import org.elasticsearch.index.reindex.UpdateByQueryRequest; import org.elasticsearch.script.Script; import org.elasticsearch.script.ScriptType; import org.springframework.util.CollectionUtils; import org.springframework.util.StringUtils; import java.io.IOException; import java.util.*; /** * Elasticsearch 文档操作服务(增强版) * * 特性: * - 完善的参数校验 * - 批量操作分批处理 * - 自动重试机制 * - 详细的操作日志 */ @Slf4j public class ElasticsearchDocumentService { private final RestHighLevelClient client; private final ElasticsearchProperties properties; public ElasticsearchDocumentService(RestHighLevelClient client, ElasticsearchProperties properties) { this.client = client; this.properties = properties; } // ==================== 新增操作 ==================== /** * 添加文档(自动生成ID) */ public String add(String indexName, Object data) { return add(indexName, null, data); } /** * 添加文档(指定ID) */ public String add(String indexName, String id, Object data) { ParamValidator.validateIndexName(indexName); ParamValidator.notNull(data, "文档数据不能为空"); return RetryUtil.executeWithRetry( () -> doAdd(indexName, id, data), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "ADD_DOCUMENT" ); } private String doAdd(String indexName, String id, Object data) throws IOException { IndexRequest request = new IndexRequest(indexName); if (StringUtils.hasText(id)) { request.id(id); } String jsonData = JSON.toJSONString(data); request.source(jsonData, XContentType.JSON); // 兼容 ES 6.x if (properties.getEnableType()) { request.type(properties.getDefaultType()); } // 设置刷新策略 setRefreshPolicy(request); IndexResponse response = client.index(request, RequestOptions.DEFAULT); log.info("添加文档成功, index: {}, id: {}", indexName, response.getId()); return response.getId(); } /** * 批量添加文档 * 自动分批处理,避免大批量数据导致的内存溢出 */ public BulkResult batchAdd(String indexName, List dataList) { ParamValidator.validateIndexName(indexName); if (CollectionUtils.isEmpty(dataList)) { return BulkResult.empty(); } int batchSize = properties.getBulkBatchSize(); int totalSize = dataList.size(); int batchCount = (int) Math.ceil((double) totalSize / batchSize); BulkResult.Builder resultBuilder = BulkResult.builder(); for (int i = 0; i < batchCount; i++) { int fromIndex = i * batchSize; int toIndex = Math.min(fromIndex + batchSize, totalSize); List batchData = dataList.subList(fromIndex, toIndex); try { BulkResult batchResult = doBatchAdd(indexName, batchData); resultBuilder.merge(batchResult); log.debug("批量添加进度: {}/{}, 成功: {}, 失败: {}", toIndex, totalSize, batchResult.getSuccessCount(), batchResult.getFailureCount()); } catch (Exception e) { log.error("批量添加第 {} 批失败", i + 1, e); resultBuilder.addFailure(batchData.size(), e.getMessage()); } } BulkResult result = resultBuilder.build(); log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}", indexName, totalSize, result.getSuccessCount(), result.getFailureCount()); return result; } private BulkResult doBatchAdd(String indexName, List dataList) throws IOException { BulkRequest bulkRequest = new BulkRequest(); for (Object data : dataList) { IndexRequest request = new IndexRequest(indexName); request.source(JSON.toJSONString(data), XContentType.JSON); if (properties.getEnableType()) { request.type(properties.getDefaultType()); } bulkRequest.add(request); } // 设置刷新策略 setRefreshPolicy(bulkRequest); BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT); return parseBulkResponse(response); } /** * 批量添加文档(带ID) */ public BulkResult batchAddWithId(String indexName, Map dataMap) { ParamValidator.validateIndexName(indexName); if (CollectionUtils.isEmpty(dataMap)) { return BulkResult.empty(); } int batchSize = properties.getBulkBatchSize(); List> entries = new ArrayList<>(dataMap.entrySet()); int totalSize = entries.size(); int batchCount = (int) Math.ceil((double) totalSize / batchSize); BulkResult.Builder resultBuilder = BulkResult.builder(); for (int i = 0; i < batchCount; i++) { int fromIndex = i * batchSize; int toIndex = Math.min(fromIndex + batchSize, totalSize); List> batchEntries = entries.subList(fromIndex, toIndex); try { BulkResult batchResult = doBatchAddWithId(indexName, batchEntries); resultBuilder.merge(batchResult); } catch (Exception e) { log.error("批量添加第 {} 批失败", i + 1, e); resultBuilder.addFailure(batchEntries.size(), e.getMessage()); } } BulkResult result = resultBuilder.build(); log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}", indexName, totalSize, result.getSuccessCount(), result.getFailureCount()); return result; } private BulkResult doBatchAddWithId(String indexName, List> entries) throws IOException { BulkRequest bulkRequest = new BulkRequest(); for (Map.Entry entry : entries) { IndexRequest request = new IndexRequest(indexName); request.id(entry.getKey()); request.source(JSON.toJSONString(entry.getValue()), XContentType.JSON); if (properties.getEnableType()) { request.type(properties.getDefaultType()); } bulkRequest.add(request); } setRefreshPolicy(bulkRequest); BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT); return parseBulkResponse(response); } // ==================== 查询操作 ==================== /** * 根据ID查询 */ public Optional getById(String indexName, String id, Class clazz) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); return RetryUtil.executeWithRetry( () -> doGetById(indexName, id, clazz), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "GET_BY_ID" ); } private Optional doGetById(String indexName, String id, Class clazz) throws IOException { GetRequest request = new GetRequest(indexName, id); GetResponse response = client.get(request, RequestOptions.DEFAULT); if (!response.isExists()) { return Optional.empty(); } String source = response.getSourceAsString(); if (source == null || source.isEmpty()) { return Optional.empty(); } T obj = JSON.parseObject(source, clazz); return Optional.of(obj); } /** * 批量根据ID查询 */ public Map getByIds(String indexName, List ids, Class clazz) { ParamValidator.validateIndexName(indexName); if (CollectionUtils.isEmpty(ids)) { return Collections.emptyMap(); } try { MultiGetRequest request = new MultiGetRequest(); for (String id : ids) { request.add(new MultiGetRequest.Item(indexName, id)); } MultiGetResponse response = client.mget(request, RequestOptions.DEFAULT); Map result = new HashMap<>(); for (MultiGetItemResponse itemResponse : response.getResponses()) { if (!itemResponse.isFailed() && itemResponse.getResponse().isExists()) { String source = itemResponse.getResponse().getSourceAsString(); if (source != null && !source.isEmpty()) { T obj = JSON.parseObject(source, clazz); result.put(itemResponse.getId(), obj); } } } return result; } catch (IOException e) { throw new ElasticsearchException("GET_BY_IDS", indexName, e); } } /** * 判断文档是否存在 */ public boolean exists(String indexName, String id) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); try { GetRequest request = new GetRequest(indexName, id); request.fetchSourceContext( org.elasticsearch.search.fetch.subphase.FetchSourceContext.DO_NOT_FETCH_SOURCE); GetResponse response = client.get(request, RequestOptions.DEFAULT); return response.isExists(); } catch (IOException e) { throw new ElasticsearchException("EXISTS", indexName, e); } } // ==================== 更新操作 ==================== /** * 更新文档(全量更新,不存在则新增) */ public boolean upsert(String indexName, String id, Object data) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); ParamValidator.notNull(data, "更新数据不能为空"); return RetryUtil.executeWithRetry( () -> doUpsert(indexName, id, data), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "UPSERT" ); } private boolean doUpsert(String indexName, String id, Object data) throws IOException { UpdateRequest request = new UpdateRequest(indexName, id); request.doc(JSON.toJSONString(data), XContentType.JSON); request.docAsUpsert(true); setRefreshPolicy(request); UpdateResponse response = client.update(request, RequestOptions.DEFAULT); log.info("更新文档成功, index: {}, id: {}, result: {}", indexName, id, response.getResult()); return response.getResult() != DocWriteResponse.Result.NOOP; } /** * 部分更新文档 */ public boolean partialUpdate(String indexName, String id, Map fields) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); if (CollectionUtils.isEmpty(fields)) { return false; } return RetryUtil.executeWithRetry( () -> doPartialUpdate(indexName, id, fields), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "PARTIAL_UPDATE" ); } private boolean doPartialUpdate(String indexName, String id, Map fields) throws IOException { UpdateRequest request = new UpdateRequest(indexName, id); request.doc(fields); setRefreshPolicy(request); UpdateResponse response = client.update(request, RequestOptions.DEFAULT); log.info("部分更新文档成功, index: {}, id: {}", indexName, id); return response.getResult() != DocWriteResponse.Result.NOOP; } /** * 使用脚本更新 */ public boolean updateByScript(String indexName, String id, String script, Map params) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); ParamValidator.notBlank(script, "脚本不能为空"); try { UpdateRequest request = new UpdateRequest(indexName, id); Script painlessScript = new Script( ScriptType.INLINE, "painless", script, params != null ? params : Collections.emptyMap() ); request.script(painlessScript); setRefreshPolicy(request); UpdateResponse response = client.update(request, RequestOptions.DEFAULT); log.info("脚本更新文档成功, index: {}, id: {}", indexName, id); return response.getResult() != DocWriteResponse.Result.NOOP; } catch (IOException e) { throw new ElasticsearchException("UPDATE_BY_SCRIPT", indexName, e); } } /** * 根据条件批量更新 */ public long updateByQuery(String indexName, QueryBuilder query, String script, Map params) { ParamValidator.validateIndexName(indexName); ParamValidator.notNull(query, "查询条件不能为空"); ParamValidator.notBlank(script, "脚本不能为空"); try { UpdateByQueryRequest request = new UpdateByQueryRequest(indexName); request.setQuery(query); request.setScript(new Script( ScriptType.INLINE, "painless", script, params != null ? params : Collections.emptyMap() )); request.setRefresh(true); BulkByScrollResponse response = client.updateByQuery( request, RequestOptions.DEFAULT); log.info("根据条件更新文档成功, index: {}, updated: {}", indexName, response.getUpdated()); return response.getUpdated(); } catch (IOException e) { throw new ElasticsearchException("UPDATE_BY_QUERY", indexName, e); } } /** * 批量更新 */ public BulkResult batchUpdate(String indexName, Map dataMap) { ParamValidator.validateIndexName(indexName); if (CollectionUtils.isEmpty(dataMap)) { return BulkResult.empty(); } try { BulkRequest bulkRequest = new BulkRequest(); dataMap.forEach((id, data) -> { UpdateRequest request = new UpdateRequest(indexName, id); request.doc(JSON.toJSONString(data), XContentType.JSON); request.docAsUpsert(true); bulkRequest.add(request); }); setRefreshPolicy(bulkRequest); BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT); BulkResult result = parseBulkResponse(response); log.info("批量更新完成, index: {}, 成功: {}, 失败: {}", indexName, result.getSuccessCount(), result.getFailureCount()); return result; } catch (IOException e) { throw new ElasticsearchException("BATCH_UPDATE", indexName, e); } } // ==================== 删除操作 ==================== /** * 删除文档 */ public boolean delete(String indexName, String id) { ParamValidator.validateIndexName(indexName); ParamValidator.validateDocumentId(id); return RetryUtil.executeWithRetry( () -> doDelete(indexName, id), properties.getRetryTimes(), properties.getRetryInterval(), RetryUtil.defaultRetryCondition(), "DELETE" ); } private boolean doDelete(String indexName, String id) throws IOException { DeleteRequest request = new DeleteRequest(indexName, id); setRefreshPolicy(request); DeleteResponse response = client.delete(request, RequestOptions.DEFAULT); log.info("删除文档成功, index: {}, id: {}", indexName, id); return response.getResult() == DocWriteResponse.Result.DELETED; } /** * 批量删除 */ public BulkResult batchDelete(String indexName, List ids) { ParamValidator.validateIndexName(indexName); if (CollectionUtils.isEmpty(ids)) { return BulkResult.empty(); } try { BulkRequest bulkRequest = new BulkRequest(); for (String id : ids) { DeleteRequest request = new DeleteRequest(indexName, id); bulkRequest.add(request); } setRefreshPolicy(bulkRequest); BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT); BulkResult result = parseBulkResponse(response); log.info("批量删除完成, index: {}, 成功: {}, 失败: {}", indexName, result.getSuccessCount(), result.getFailureCount()); return result; } catch (IOException e) { throw new ElasticsearchException("BATCH_DELETE", indexName, e); } } /** * 根据条件删除 */ public long deleteByQuery(String indexName, QueryBuilder query) { ParamValidator.validateIndexName(indexName); ParamValidator.notNull(query, "查询条件不能为空"); try { DeleteByQueryRequest request = new DeleteByQueryRequest(indexName); request.setQuery(query); request.setRefresh(true); BulkByScrollResponse response = client.deleteByQuery( request, RequestOptions.DEFAULT); log.info("根据条件删除文档成功, index: {}, deleted: {}", indexName, response.getDeleted()); return response.getDeleted(); } catch (IOException e) { throw new ElasticsearchException("DELETE_BY_QUERY", indexName, e); } } // ==================== 内部工具方法 ==================== /** * 设置刷新策略 */ private void setRefreshPolicy(IndexRequest request) { String policy = properties.getBulkRefreshPolicy(); if ("immediate".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); } else if ("wait_for".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL); } } private void setRefreshPolicy(UpdateRequest request) { String policy = properties.getBulkRefreshPolicy(); if ("immediate".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); } else if ("wait_for".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL); } } private void setRefreshPolicy(DeleteRequest request) { String policy = properties.getBulkRefreshPolicy(); if ("immediate".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); } else if ("wait_for".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL); } } private void setRefreshPolicy(BulkRequest request) { String policy = properties.getBulkRefreshPolicy(); if ("immediate".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE); } else if ("wait_for".equals(policy)) { request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL); } } /** * 解析批量操作响应 */ private BulkResult parseBulkResponse(BulkResponse response) { BulkResult.Builder builder = BulkResult.builder(); for (BulkItemResponse itemResponse : response.getItems()) { if (itemResponse.isFailed()) { builder.addFailure(itemResponse.getId(), itemResponse.getFailureMessage()); } else { builder.addSuccess(itemResponse.getId()); } } return builder.build(); } // ==================== 批量操作结果 ==================== /** * 批量操作结果 */ @lombok.Data public static class BulkResult { private int successCount; private int failureCount; private List successIds; private List failures; @lombok.Data @lombok.AllArgsConstructor public static class FailureItem { private String id; private String reason; } public static BulkResult empty() { BulkResult result = new BulkResult(); result.successCount = 0; result.failureCount = 0; result.successIds = Collections.emptyList(); result.failures = Collections.emptyList(); return result; } public boolean hasFailures() { return failureCount > 0; } public boolean isAllSuccess() { return failureCount == 0; } public static Builder builder() { return new Builder(); } public static class Builder { private List successIds = new ArrayList<>(); private List failures = new ArrayList<>(); public Builder addSuccess(String id) { successIds.add(id); return this; } public Builder addFailure(String id, String reason) { failures.add(new FailureItem(id, reason)); return this; } public Builder addFailure(int count, String reason) { for (int i = 0; i < count; i++) { failures.add(new FailureItem(null, reason)); } return this; } public Builder merge(BulkResult other) { if (other.successIds != null) { successIds.addAll(other.successIds); } if (other.failures != null) { failures.addAll(other.failures); } return this; } public BulkResult build() { BulkResult result = new BulkResult(); result.successIds = successIds; result.failures = failures; result.successCount = successIds.size(); result.failureCount = failures.size(); return result; } } } } ``` ### **9.12 滚动查询结果类** ```java package com.example.elasticsearch.result; import lombok.AllArgsConstructor; import lombok.Data; import lombok.NoArgsConstructor; import java.util.List; /** * 滚动查询结果 */ @Data @NoArgsConstructor @AllArgsConstructor public class ScrollResult { /** * 滚动ID(用于获取下一批数据) */ private String scrollId; /** * 当前批次数据 */ private List records; /** * 总记录数 */ private Long total; /** * 是否是最后一批(没有更多数据) */ private Boolean finished; /** * 判断是否还有更多数据 */ public boolean hasMore() { return !finished && records != null && !records.isEmpty(); } } ``` ### **9.13 索引操作服务** ```java package com.example.elasticsearch.core; import com.example.elasticsearch.config.ElasticsearchProperties; import com.example.elasticsearch.exception.ElasticsearchException; import lombok.extern.slf4j.Slf4j; import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest; import org.elasticsearch.action.support.master.AcknowledgedResponse; import org.elasticsearch.client.RequestOptions; import org.elasticsearch.client.RestHighLevelClient; import org.elasticsearch.client.indices.*; import org.elasticsearch.common.settings.Settings; import org.elasticsearch.common.xcontent.XContentType; import java.io.IOException; import java.util.Map; /** * Elasticsearch 索引操作服务 */ @Slf4j public class ElasticsearchIndexService { private final RestHighLevelClient client; private final ElasticsearchProperties properties; public ElasticsearchIndexService(RestHighLevelClient client, ElasticsearchProperties properties) { this.client = client; this.properties = properties; } /** * 创建索引 */ public boolean createIndex(String indexName) { return createIndex(indexName, null, null); } /** * 创建索引(带 Mapping) */ public boolean createIndex(String indexName, String mapping) { try { CreateIndexRequest request = new CreateIndexRequest(indexName); if (mapping != null) { request.source(mapping, XContentType.JSON); } CreateIndexResponse response = client.indices() .create(request, RequestOptions.DEFAULT); log.info("创建索引成功, index: {}", indexName); return response.isAcknowledged(); } catch (IOException e) { log.error("创建索引失败", e); throw new ElasticsearchException("CREATE_INDEX", indexName, e); } } /** * 创建索引(带 Settings 和 Mapping) */ public boolean createIndex(String indexName, Map settings, Map mapping) { try { CreateIndexRequest request = new CreateIndexRequest(indexName); if (settings != null) { request.settings(settings); } if (mapping != null) { request.mapping(mapping); } CreateIndexResponse response = client.indices() .create(request, RequestOptions.DEFAULT); log.info("创建索引成功, index: {}", indexName); return response.isAcknowledged(); } catch (IOException e) { log.error("创建索引失败", e); throw new ElasticsearchException("CREATE_INDEX", indexName, e); } } /** * 判断索引是否存在 */ public boolean existsIndex(String indexName) { try { GetIndexRequest request = new GetIndexRequest(indexName); return client.indices().exists(request, RequestOptions.DEFAULT); } catch (IOException e) { log.error("判断索引是否存在失败", e); throw new ElasticsearchException("EXISTS_INDEX", indexName, e); } } /** * 删除索引 */ public boolean deleteIndex(String indexName) { try { DeleteIndexRequest request = new DeleteIndexRequest(indexName); AcknowledgedResponse response = client.indices() .delete(request, RequestOptions.DEFAULT); log.info("删除索引成功, index: {}", indexName); return response.isAcknowledged(); } catch (IOException e) { log.error("删除索引失败", e); throw new ElasticsearchException("DELETE_INDEX", indexName, e); } } /** * 更新 Mapping */ public boolean updateMapping(String indexName, Map properties) { try { PutMappingRequest request = new PutMappingRequest(indexName); request.source(Map.of("properties", properties)); AcknowledgedResponse response = client.indices() .putMapping(request, RequestOptions.DEFAULT); log.info("更新 Mapping 成功, index: {}", indexName); return response.isAcknowledged(); } catch (IOException e) { log.error("更新 Mapping 失败", e); throw new ElasticsearchException("UPDATE_MAPPING", indexName, e); } } /** * 获取索引信息 */ public GetIndexResponse getIndex(String indexName) { try { GetIndexRequest request = new GetIndexRequest(indexName); return client.indices().get(request, RequestOptions.DEFAULT); } catch (IOException e) { log.error("获取索引信息失败", e); throw new ElasticsearchException("GET_INDEX", indexName, e); } } /** * 刷新索引 */ public void refreshIndex(String... indexNames) { try { org.elasticsearch.client.indices.RefreshRequest request = new org.elasticsearch.client.indices.RefreshRequest(indexNames); client.indices().refresh(request, RequestOptions.DEFAULT); log.info("刷新索引成功, indices: {}", String.join(",", indexNames)); } catch (IOException e) { log.error("刷新索引失败", e); throw new ElasticsearchException("REFRESH_INDEX", String.join(",", indexNames), e); } } /** * 创建索引别名 */ public boolean createAlias(String indexName, String aliasName) { try { var request = new org.elasticsearch.client.indices.IndicesAliasesRequest(); var action = new org.elasticsearch.client.indices.IndicesAliasesRequest .AliasActions(org.elasticsearch.client.indices.IndicesAliasesRequest .AliasActions.Type.ADD) .index(indexName) .alias(aliasName); request.addAliasAction(action); var response = client.indices().updateAliases(request, RequestOptions.DEFAULT); log.info("创建索引别名成功, index: {}, alias: {}", indexName, aliasName); return response.isAcknowledged(); } catch (IOException e) { log.error("创建索引别名失败", e); throw new ElasticsearchException("CREATE_ALIAS", indexName, e); } } } ``` ### **9.14 架构图** ![[7-Blog/后端与微服务/assets/架构图-39a2cd40.jpg]] 1. ✅ **链式调用** - `EsQueryBuilder` 支持流畅的链式构建 2. ✅ **类型安全** - 泛型支持,自动转换结果类型 3. ✅ **功能完整** - 覆盖索引、文档、查询、聚合操作 4. ✅ **高亮支持** - 内置高亮配置和结果解析 5. ✅ **分页封装** - 统一的分页结果对象 6. ✅ **连接池** - 支持集群、认证、连接池配置 7. ✅ **版本兼容** - 支持 ES 6.x/7.x type 兼容 8. ✅ **异常处理** - 统一的异常封装 ## **十、完整使用文档** ### **快速开始** #### **添加依赖** ``` com.example elasticsearch-spring-boot-starter 1.0.0 org.springframework.boot spring-boot-starter-data-elasticsearch com.alibaba fastjson 1.2.83 ``` #### **配置文件** ``` # application.yml elasticsearch: enabled: true nodes: - localhost:9200 username: elastic # 可选 password: your_password # 可选 scheme: http # 连接配置 connect-timeout: 5000 socket-timeout: 30000 max-conn-total: 100 max-conn-per-route: 50 # 重试配置 retry-times: 3 retry-interval: 1000 # 批量操作配置 bulk-batch-size: 1000 bulk-refresh-policy: none # none, immediate, wait_for # 日志配置 print-dsl: true # 开发环境建议开启 slow-query-threshold: 3000 # 慢查询阈值(ms) # 版本兼容 enable-type: false # ES 7.x 设为 false ``` #### **创建实体类** ```java package com.example.demo.entity; import lombok.Data; import java.util.Date; /** * 节目实体 */ @Data public class Program { private Long id; /** * 节目标题 */ private String title; /** * 演员 */ private String actor; /** * 分类ID */ private Long categoryId; /** * 分类名称 */ private String categoryName; /** * 地区ID */ private Long areaId; /** * 地区名称 */ private String areaName; /** * 演出时间 */ private Date showTime; /** * 最低价格 */ private Double minPrice; /** * 最高价格 */ private Double maxPrice; /** * 状态:1-上架, 0-下架 */ private Integer status; /** * 封面图片 */ private String coverImage; /** * 创建时间 */ private Date createTime; /** * 更新时间 */ private Date updateTime; } ``` #### **创建搜索参数类** ```java package com.example.demo.dto; import lombok.Data; import java.util.Date; /** * 节目搜索参数 */ @Data public class ProgramSearchParam { /** * 关键词(搜索标题和演员) */ private String keyword; /** * 分类ID */ private Long categoryId; /** * 地区ID */ private Long areaId; /** * 状态 */ private Integer status; /** * 最低价格 */ private Double minPrice; /** * 最高价格 */ private Double maxPrice; /** * 开始时间 */ private Date startTime; /** * 结束时间 */ private Date endTime; /** * 页码 */ private Integer pageNum = 1; /** * 每页大小 */ private Integer pageSize = 10; /** * 排序字段 */ private String sortField = "showTime"; /** * 排序方式:asc, desc */ private String sortOrder = "asc"; } ``` ### **完整业务示例** ```java package com.example.demo.service; import com.example.demo.dto.ProgramSearchParam; import com.example.demo.entity.Program; import com.example.elasticsearch.core.ElasticsearchDocumentService; import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult; import com.example.elasticsearch.core.ElasticsearchIndexService; import com.example.elasticsearch.core.ElasticsearchService; import com.example.elasticsearch.query.EsQueryBuilder; import com.example.elasticsearch.result.AggregationResult; import com.example.elasticsearch.result.PageResult; import com.example.elasticsearch.result.ScrollResult; import com.example.elasticsearch.result.SearchResult; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; import org.elasticsearch.search.sort.SortOrder; import org.springframework.stereotype.Service; import javax.annotation.PostConstruct; import java.util.*; import java.util.function.Consumer; import java.util.stream.Collectors; /** * 节目搜索服务 - 完整示例 */ @Slf4j @Service @RequiredArgsConstructor public class ProgramSearchService { private final ElasticsearchService esService; private final ElasticsearchDocumentService documentService; private final ElasticsearchIndexService indexService; private static final String INDEX_NAME = "program"; // ==================== 索引管理 ==================== /** * 初始化索引(应用启动时调用) */ @PostConstruct public void initIndex() { if (!indexService.existsIndex(INDEX_NAME)) { String mapping = """ { "mappings": { "properties": { "id": { "type": "long" }, "title": { "type": "text", "analyzer": "ik_max_word", "search_analyzer": "ik_smart", "fields": { "keyword": { "type": "keyword" } } }, "actor": { "type": "text", "analyzer": "ik_max_word", "fields": { "keyword": { "type": "keyword" } } }, "categoryId": { "type": "long" }, "categoryName": { "type": "keyword" }, "areaId": { "type": "long" }, "areaName": { "type": "keyword" }, "showTime": { "type": "date" }, "minPrice": { "type": "double" }, "maxPrice": { "type": "double" }, "status": { "type": "integer" }, "coverImage": { "type": "keyword", "index": false }, "createTime": { "type": "date" }, "updateTime": { "type": "date" } } }, "settings": { "number_of_shards": 3, "number_of_replicas": 1, "refresh_interval": "1s" } } """; indexService.createIndex(INDEX_NAME, mapping); log.info("索引 [{}] 创建成功", INDEX_NAME); } } /** * 重建索引 */ public void rebuildIndex() { // 删除旧索引 if (indexService.existsIndex(INDEX_NAME)) { indexService.deleteIndex(INDEX_NAME); } // 重新初始化 initIndex(); } // ==================== 文档操作 ==================== /** * 添加/更新节目 */ public void saveProgram(Program program) { program.setUpdateTime(new Date()); if (program.getCreateTime() == null) { program.setCreateTime(new Date()); } documentService.upsert(INDEX_NAME, String.valueOf(program.getId()), program); } /** * 批量添加节目 */ public BulkResult batchSavePrograms(List programs) { Date now = new Date(); programs.forEach(p -> { p.setUpdateTime(now); if (p.getCreateTime() == null) { p.setCreateTime(now); } }); return documentService.batchAdd(INDEX_NAME, programs); } /** * 批量添加节目(带ID) */ public BulkResult batchSaveProgramsWithId(List programs) { Date now = new Date(); Map dataMap = new HashMap<>(); programs.forEach(p -> { p.setUpdateTime(now); if (p.getCreateTime() == null) { p.setCreateTime(now); } dataMap.put(String.valueOf(p.getId()), p); }); return documentService.batchAddWithId(INDEX_NAME, dataMap); } /** * 根据ID查询节目 */ public Optional getProgramById(Long id) { return documentService.getById(INDEX_NAME, String.valueOf(id), Program.class); } /** * 批量查询节目 */ public Map getProgramsByIds(List ids) { List strIds = ids.stream() .map(String::valueOf) .collect(Collectors.toList()); Map resultMap = documentService.getByIds( INDEX_NAME, strIds, Program.class); // 转换 key 类型 return resultMap.entrySet().stream() .collect(Collectors.toMap( e -> Long.parseLong(e.getKey()), Map.Entry::getValue )); } /** * 删除节目 */ public void deleteProgram(Long id) { documentService.delete(INDEX_NAME, String.valueOf(id)); } /** * 批量删除节目 */ public BulkResult batchDeletePrograms(List ids) { List strIds = ids.stream() .map(String::valueOf) .collect(Collectors.toList()); return documentService.batchDelete(INDEX_NAME, strIds); } /** * 更新节目状态 */ public boolean updateProgramStatus(Long id, Integer status) { Map fields = new HashMap<>(); fields.put("status", status); fields.put("updateTime", new Date()); return documentService.partialUpdate(INDEX_NAME, String.valueOf(id), fields); } /** * 批量更新节目状态(使用脚本) */ public long batchUpdateStatus(List ids, Integer status) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .filterTerms("id", ids); String script = "ctx._source.status = params.status; " + "ctx._source.updateTime = params.updateTime"; Map params = new HashMap<>(); params.put("status", status); params.put("updateTime", System.currentTimeMillis()); return documentService.updateByQuery( INDEX_NAME, queryBuilder.getBoolQuery(), script, params); } // ==================== 基础搜索 ==================== /** * 根据分类查询节目列表 */ public List findByCategory(Long categoryId) { return esService.query(INDEX_NAME, "categoryId", categoryId, Program.class); } /** * 根据多条件查询 */ public List findByConditions(Long categoryId, Long areaId, Integer status) { Map params = new HashMap<>(); params.put("categoryId", categoryId); params.put("areaId", areaId); params.put("status", status); return esService.query(INDEX_NAME, params, Program.class); } /** * 简单分页查询 */ public PageResult findPage(int pageNum, int pageSize) { return esService.queryPage(INDEX_NAME, null, pageNum, pageSize, Program.class); } // ==================== 高级搜索 ==================== /** * 关键词搜索(标题 + 演员) */ public List searchByKeyword(String keyword) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(keyword, "title", "actor") .filterTerm("status", 1) // 只搜索上架的 .sort("showTime", SortOrder.ASC); return esService.search(queryBuilder, Program.class); } /** * 关键词搜索(带高亮) */ public List> searchWithHighlight(String keyword) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(keyword, "title", "actor") .filterTerm("status", 1) .highlight("", "", "title", "actor") .sort("showTime", SortOrder.ASC); return esService.searchWithHighlight(queryBuilder, Program.class); } /** * 复杂条件搜索 + 分页 */ public PageResult search(ProgramSearchParam param) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) // 关键词搜索(标题或演员) .shouldMatch(param.getKeyword(), "title", "actor") // 分类过滤 .filterTerm("categoryId", param.getCategoryId()) // 地区过滤 .filterTerm("areaId", param.getAreaId()) // 状态过滤 .filterTerm("status", param.getStatus()) // 价格范围 .filterRange("minPrice", param.getMinPrice(), param.getMaxPrice()) // 时间范围 .filterRange("showTime", param.getStartTime() != null ? param.getStartTime().getTime() : null, param.getEndTime() != null ? param.getEndTime().getTime() : null) // 排序 .sort(param.getSortField(), "desc".equalsIgnoreCase(param.getSortOrder()) ? SortOrder.DESC : SortOrder.ASC); return esService.searchPage(queryBuilder, param.getPageNum(), param.getPageSize(), Program.class); } /** * 复杂条件搜索 + 分页 + 高亮 */ public PageResult> searchWithHighlight(ProgramSearchParam param) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(param.getKeyword(), "title", "actor") .filterTerm("categoryId", param.getCategoryId()) .filterTerm("areaId", param.getAreaId()) .filterTerm("status", param.getStatus()) .filterRange("minPrice", param.getMinPrice(), param.getMaxPrice()) .filterRange("showTime", param.getStartTime() != null ? param.getStartTime().getTime() : null, param.getEndTime() != null ? param.getEndTime().getTime() : null) .highlight("title", "actor") .sort(param.getSortField(), "desc".equalsIgnoreCase(param.getSortOrder()) ? SortOrder.DESC : SortOrder.ASC); return esService.searchPageWithHighlight(queryBuilder, param.getPageNum(), param.getPageSize(), Program.class); } // ==================== 深度分页(Search After) ==================== /** * 使用 Search After 进行深度分页 * 适用于页码超过 1000 的场景 * * @param param 搜索参数 * @param searchAfterValue 上一页最后一条记录的排序值 * @return 搜索结果(包含 sortValues 用于下一页查询) */ public List> searchAfter(ProgramSearchParam param, Object[] searchAfterValue) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(param.getKeyword(), "title", "actor") .filterTerm("categoryId", param.getCategoryId()) .filterTerm("status", param.getStatus()) // 必须有排序字段,且最后一个字段应该是唯一的(如 _id) .sort(param.getSortField(), "desc".equalsIgnoreCase(param.getSortOrder()) ? SortOrder.DESC : SortOrder.ASC) .sort("id", SortOrder.ASC); // 使用 id 保证唯一性 return esService.searchAfter(queryBuilder, searchAfterValue, param.getPageSize(), Program.class); } // ==================== 滚动查询(大数据导出) ==================== /** * 滚动查询导出所有数据 * 适用于数据导出场景 * * @param param 搜索参数 * @param consumer 数据消费者 */ public void scrollExport(ProgramSearchParam param, Consumer> consumer) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(param.getKeyword(), "title", "actor") .filterTerm("categoryId", param.getCategoryId()) .filterTerm("status", param.getStatus()) .sort("id", SortOrder.ASC); // 滚动查询必须有排序 int batchSize = 1000; int scrollTime = 5; // 滚动上下文保持时间(分钟) // 初始化滚动查询 ScrollResult scrollResult = esService.scrollSearch( queryBuilder, scrollTime, batchSize, Program.class); try { // 处理首批数据 if (!scrollResult.getRecords().isEmpty()) { consumer.accept(scrollResult.getRecords()); } // 继续获取后续数据 while (scrollResult.hasMore()) { scrollResult = esService.scrollNext( scrollResult.getScrollId(), scrollTime, batchSize, Program.class); if (!scrollResult.getRecords().isEmpty()) { consumer.accept(scrollResult.getRecords()); } } log.info("滚动导出完成, 总数: {}", scrollResult.getTotal()); } finally { // 清除滚动上下文 esService.clearScroll(scrollResult.getScrollId()); } } // ==================== 聚合统计 ==================== /** * 按分类统计节目数量 */ public Map countByCategory() { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .filterTerm("status", 1) // 只统计上架的 .termsAggregation("category_count", "categoryId", 100); Map aggResult = esService.searchAggregation(queryBuilder); AggregationResult categoryAgg = aggResult.get("category_count"); if (categoryAgg == null || categoryAgg.getBuckets() == null) { return Collections.emptyMap(); } return categoryAgg.getBuckets().stream() .collect(Collectors.toMap( AggregationResult.BucketData::getKey, AggregationResult.BucketData::getDocCount )); } /** * 统计价格区间 */ public Map getPriceStats() { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .filterTerm("status", 1) .minAggregation("min_price", "minPrice") .maxAggregation("max_price", "maxPrice") .avgAggregation("avg_price", "minPrice"); Map aggResult = esService.searchAggregation(queryBuilder); Map stats = new HashMap<>(); if (aggResult.containsKey("min_price")) { stats.put("minPrice", aggResult.get("min_price").getValue()); } if (aggResult.containsKey("max_price")) { stats.put("maxPrice", aggResult.get("max_price").getValue()); } if (aggResult.containsKey("avg_price")) { stats.put("avgPrice", aggResult.get("avg_price").getValue()); } return stats; } /** * 综合统计 */ public Map getStatistics() { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .filterTerm("status", 1) .termsAggregation("by_category", "categoryId", 50) .termsAggregation("by_area", "areaId", 50) .avgAggregation("avg_price", "minPrice") .minAggregation("min_price", "minPrice") .maxAggregation("max_price", "maxPrice"); Map aggResult = esService.searchAggregation(queryBuilder); Map stats = new HashMap<>(); stats.put("aggregations", aggResult); stats.put("total", esService.count( EsQueryBuilder.builder(INDEX_NAME).filterTerm("status", 1))); return stats; } // ==================== 工具方法 ==================== /** * 统计符合条件的数量 */ public long count(ProgramSearchParam param) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .shouldMatch(param.getKeyword(), "title", "actor") .filterTerm("categoryId", param.getCategoryId()) .filterTerm("areaId", param.getAreaId()) .filterTerm("status", param.getStatus()); return esService.count(queryBuilder); } /** * 判断是否存在符合条件的数据 */ public boolean exists(ProgramSearchParam param) { return count(param) > 0; } } ``` ### **Controller 示例** ```java package com.example.demo.controller; import com.example.demo.dto.ProgramSearchParam; import com.example.demo.entity.Program; import com.example.demo.service.ProgramSearchService; import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult; import com.example.elasticsearch.result.PageResult; import com.example.elasticsearch.result.SearchResult; import lombok.RequiredArgsConstructor; import org.springframework.web.bind.annotation.*; import java.util.List; import java.util.Map; import java.util.Optional; /** * 节目搜索 API */ @RestController @RequestMapping("/api/programs") @RequiredArgsConstructor public class ProgramController { private final ProgramSearchService searchService; // ==================== 基础 CRUD ==================== /** * 添加/更新节目 */ @PostMapping public String saveProgram(@RequestBody Program program) { searchService.saveProgram(program); return "success"; } /** * 批量添加节目 */ @PostMapping("/batch") public BulkResult batchSavePrograms(@RequestBody List programs) { return searchService.batchSavePrograms(programs); } /** * 根据ID查询 */ @GetMapping("/{id}") public Optional getProgram(@PathVariable Long id) { return searchService.getProgramById(id); } /** * 批量查询 */ @PostMapping("/batch-get") public Map batchGetPrograms(@RequestBody List ids) { return searchService.getProgramsByIds(ids); } /** * 删除节目 */ @DeleteMapping("/{id}") public String deleteProgram(@PathVariable Long id) { searchService.deleteProgram(id); return "success"; } /** * 更新状态 */ @PutMapping("/{id}/status") public boolean updateStatus(@PathVariable Long id, @RequestParam Integer status) { return searchService.updateProgramStatus(id, status); } // ==================== 搜索 ==================== /** * 关键词搜索 */ @GetMapping("/search") public List search(@RequestParam String keyword) { return searchService.searchByKeyword(keyword); } /** * 关键词搜索(带高亮) */ @GetMapping("/search/highlight") public List> searchWithHighlight(@RequestParam String keyword) { return searchService.searchWithHighlight(keyword); } /** * 高级搜索(分页) */ @PostMapping("/search/page") public PageResult searchPage(@RequestBody ProgramSearchParam param) { return searchService.search(param); } /** * 高级搜索(分页 + 高亮) */ @PostMapping("/search/page-highlight") public PageResult> searchPageWithHighlight( @RequestBody ProgramSearchParam param) { return searchService.searchWithHighlight(param); } /** * 深度分页(Search After) */ @PostMapping("/search/after") public List> searchAfter( @RequestBody ProgramSearchParam param, @RequestParam(required = false) Long afterId, @RequestParam(required = false) Long afterShowTime) { Object[] searchAfter = null; if (afterShowTime != null && afterId != null) { searchAfter = new Object[]{afterShowTime, afterId}; } return searchService.searchAfter(param, searchAfter); } // ==================== 统计 ==================== /** * 按分类统计 */ @GetMapping("/stats/category") public Map countByCategory() { return searchService.countByCategory(); } /** * 价格统计 */ @GetMapping("/stats/price") public Map getPriceStats() { return searchService.getPriceStats(); } /** * 综合统计 */ @GetMapping("/stats") public Map getStatistics() { return searchService.getStatistics(); } // ==================== 索引管理 ==================== /** * 重建索引 */ @PostMapping("/index/rebuild") public String rebuildIndex() { searchService.rebuildIndex(); return "success"; } } ``` ### **最佳实践指南** #### **索引设计** ```java /** * 索引设计最佳实践 */ public class IndexDesignBestPractice { /** * 1. 合理设置分片数 * - 单个分片大小建议 10-50GB * - 分片数 = 数据量 / 单分片大小 */ public static final int SHARD_COUNT = 3; /** * 2. 副本数设置 * - 生产环境至少 1 个副本 * - 写入密集型可先设为 0,完成后再设为 1 */ public static final int REPLICA_COUNT = 1; /** * 3. Mapping 设计原则 * - 明确字段类型,避免动态映射 * - text 用于全文搜索,keyword 用于精确匹配/聚合 * - 不需要搜索的字段设置 index: false * - 大文本字段考虑使用 store: true */ public static String createOptimizedMapping() { return """ { "mappings": { "properties": { "id": { "type": "long" }, "title": { "type": "text", "analyzer": "ik_max_word", "search_analyzer": "ik_smart", "fields": { "keyword": { "type": "keyword", "ignore_above": 256 } } }, "description": { "type": "text", "analyzer": "ik_max_word", "index": true, "store": true }, "status": { "type": "keyword" }, "price": { "type": "scaled_float", "scaling_factor": 100 }, "createTime": { "type": "date", "format": "epoch_millis" }, "tags": { "type": "keyword" }, "coverImage": { "type": "keyword", "index": false } } }, "settings": { "number_of_shards": 3, "number_of_replicas": 1, "refresh_interval": "1s", "max_result_window": 10000, "analysis": { "analyzer": { "ik_smart_pinyin": { "type": "custom", "tokenizer": "ik_smart", "filter": ["lowercase", "pinyin_filter"] } } } } } """; } } ``` #### **查询优化** ```java /** * 查询优化最佳实践 */ @Slf4j @Service public class QueryOptimizationService { @Autowired private ElasticsearchService esService; private static final String INDEX_NAME = "program"; /** * 1. 使用 filter 而非 must 进行过滤 * filter 不计算评分,性能更好,且结果可缓存 */ public List optimizedFilterQuery(Long categoryId, Integer status) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) // ✅ 正确:使用 filter 进行精确匹配 .filterTerm("categoryId", categoryId) .filterTerm("status", status); // ❌ 错误:使用 must 进行精确匹配会计算评分,浪费性能 // .term("categoryId", categoryId) // .term("status", status) return esService.search(queryBuilder, Program.class); } /** * 2. 合理使用分页,避免深度分页 * from + size 超过 10000 会报错 */ public PageResult safePageQuery(int pageNum, int pageSize) { // 计算深度 int depth = (pageNum - 1) * pageSize + pageSize; if (depth > 10000) { // 深度分页使用 search_after log.warn("分页深度超过10000,建议使用 searchAfter"); throw new IllegalArgumentException("分页深度超出限制,请使用 searchAfter"); } EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .page(pageNum, pageSize) .sort("id", SortOrder.ASC); return esService.searchPage(queryBuilder, pageNum, pageSize, Program.class); } /** * 3. 只返回需要的字段,减少网络传输 */ public List selectiveFieldsQuery(String keyword) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .match("title", keyword) // 只返回需要的字段 .includes("id", "title", "minPrice", "showTime") // 排除大字段 .excludes("description", "coverImage"); return esService.search(queryBuilder, Program.class); } /** * 4. 使用 bool 查询组合多个条件 */ public List complexBoolQuery(String keyword, List categoryIds, Double minPrice, Double maxPrice, List excludeStatus) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) // must: 必须匹配(影响评分) .match("title", keyword) // filter: 必须匹配(不影响评分,可缓存) .filterTerms("categoryId", categoryIds) .filterRange("minPrice", minPrice, maxPrice) // must_not: 必须不匹配 .mustNotTerms("status", excludeStatus) // should: 可选匹配(有则加分) .should(QueryBuilders.matchQuery("actor", keyword)) .minimumShouldMatch(0); return esService.search(queryBuilder, Program.class); } /** * 5. 使用 constant_score 包装不需要评分的查询 */ public List constantScoreQuery(Long categoryId) { // 使用原生 QueryBuilder 实现 constant_score QueryBuilder query = QueryBuilders.constantScoreQuery( QueryBuilders.termQuery("categoryId", categoryId) ).boost(1.0f); return esService.searchByQueryBuilder(INDEX_NAME, query, 100, Program.class); } /** * 6. 批量查询优化:使用 multi_search */ public Map> multiSearch(List categoryIds) { // 注意:需要扩展 ElasticsearchService 支持 multi_search // 这里展示思路 Map> result = new HashMap<>(); for (Long categoryId : categoryIds) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .filterTerm("categoryId", categoryId) .fromSize(0, 10); List programs = esService.search(queryBuilder, Program.class); result.put(String.valueOf(categoryId), programs); } return result; } /** * 7. 高亮优化:限制高亮片段 */ public List> optimizedHighlightQuery(String keyword) { EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME) .match("title", keyword) .match("description", keyword); // 自定义高亮配置 HighlightBuilder highlightBuilder = new HighlightBuilder() .field(new HighlightBuilder.Field("title") .fragmentSize(100) // 片段大小 .numOfFragments(1)) // 片段数量 .field(new HighlightBuilder.Field("description") .fragmentSize(150) .numOfFragments(2)) .preTags("") .postTags(""); queryBuilder.setHighlightBuilder(highlightBuilder); return esService.searchWithHighlight(queryBuilder, Program.class); } } ``` #### **写入优化** ```java /** * 写入优化最佳实践 */ @Slf4j @Service public class WriteOptimizationService { @Autowired private ElasticsearchDocumentService documentService; @Autowired private ElasticsearchIndexService indexService; private static final String INDEX_NAME = "program"; /** * 1. 批量写入优化 * - 合理设置批次大小(通常 1000-5000) * - 控制并发写入线程数 */ public void optimizedBatchWrite(List programs) { int batchSize = 1000; int totalSize = programs.size(); log.info("开始批量写入,总数: {}", totalSize); long startTime = System.currentTimeMillis(); // 分批处理 for (int i = 0; i < totalSize; i += batchSize) { int endIndex = Math.min(i + batchSize, totalSize); List batch = programs.subList(i, endIndex); BulkResult result = documentService.batchAdd(INDEX_NAME, batch); if (result.hasFailures()) { log.warn("批次 {}-{} 部分失败,失败数: {}", i, endIndex, result.getFailureCount()); } log.info("写入进度: {}/{}", endIndex, totalSize); } long elapsed = System.currentTimeMillis() - startTime; log.info("批量写入完成,耗时: {}ms,速率: {} docs/s", elapsed, totalSize * 1000L / elapsed); } /** * 2. 大批量数据导入优化 * - 临时关闭副本 * - 增大 refresh_interval * - 导入完成后恢复 */ public void bulkImportWithOptimization(List programs) { try { // 优化索引设置 optimizeIndexForBulkImport(); // 执行批量导入 optimizedBatchWrite(programs); } finally { // 恢复索引设置 restoreIndexSettings(); } } private void optimizeIndexForBulkImport() { Map settings = new HashMap<>(); settings.put("index.refresh_interval", "-1"); // 关闭自动刷新 settings.put("index.number_of_replicas", 0); // 临时关闭副本 indexService.updateSettings(INDEX_NAME, settings); log.info("索引设置已优化"); } private void restoreIndexSettings() { Map settings = new HashMap<>(); settings.put("index.refresh_interval", "1s"); // 恢复刷新间隔 settings.put("index.number_of_replicas", 1); // 恢复副本 indexService.updateSettings(INDEX_NAME, settings); // 强制刷新 indexService.refreshIndex(INDEX_NAME); // 强制合并段 indexService.forceMerge(INDEX_NAME, 1); log.info("索引设置已恢复"); } /** * 3. 使用异步写入 */ @Async public CompletableFuture asyncBatchWrite(List programs) { BulkResult result = documentService.batchAdd(INDEX_NAME, programs); return CompletableFuture.completedFuture(result); } /** * 4. 并行批量写入 */ public void parallelBatchWrite(List programs, int threadCount) { int batchSize = programs.size() / threadCount; ExecutorService executor = Executors.newFixedThreadPool(threadCount); List> futures = new ArrayList<>(); try { for (int i = 0; i < threadCount; i++) { int start = i * batchSize; int end = (i == threadCount - 1) ? programs.size() : start + batchSize; List batch = programs.subList(start, end); Future future = executor.submit(() -> documentService.batchAdd(INDEX_NAME, batch)); futures.add(future); } // 等待所有任务完成 int totalSuccess = 0; int totalFailure = 0; for (Future future : futures) { BulkResult result = future.get(); totalSuccess += result.getSuccessCount(); totalFailure += result.getFailureCount(); } log.info("并行写入完成,成功: {}, 失败: {}", totalSuccess, totalFailure); } catch (Exception e) { log.error("并行写入失败", e); throw new RuntimeException(e); } finally { executor.shutdown(); } } } ``` #### **缓存策略** ```java /** * ES 查询缓存策略 */ @Slf4j @Service public class EsCacheService { @Autowired private ElasticsearchService esService; @Autowired private RedisTemplate redisTemplate; private static final String CACHE_PREFIX = "es:program:"; private static final Duration CACHE_TTL = Duration.ofMinutes(5); /** * 1. 使用 Redis 缓存查询结果 */ @SuppressWarnings("unchecked") public List searchWithCache(String keyword) { String cacheKey = CACHE_PREFIX + "search:" + DigestUtils.md5Hex(keyword); // 先查缓存 Object cached = redisTemplate.opsForValue().get(cacheKey); if (cached != null) { log.debug("命中缓存: {}", cacheKey); return (List) cached; } // 查询 ES EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program") .match("title", keyword); List result = esService.search(queryBuilder, Program.class); // 写入缓存 if (!result.isEmpty()) { redisTemplate.opsForValue().set(cacheKey, result, CACHE_TTL); } return result; } /** * 2. 热门搜索词预热缓存 */ @Scheduled(fixedRate = 300000) // 每5分钟 public void warmUpCache() { List hotKeywords = getHotKeywords(); for (String keyword : hotKeywords) { try { searchWithCache(keyword); } catch (Exception e) { log.warn("预热缓存失败: {}", keyword, e); } } log.info("缓存预热完成,预热 {} 个关键词", hotKeywords.size()); } /** * 3. 数据更新时清除相关缓存 */ public void invalidateCache(Long programId) { // 清除该节目相关的所有缓存 Set keys = redisTemplate.keys(CACHE_PREFIX + "*"); if (keys != null && !keys.isEmpty()) { redisTemplate.delete(keys); log.info("清除 {} 个缓存", keys.size()); } } /** * 4. 分类统计结果缓存(变化不频繁的数据) */ @Cacheable(value = "es:stats", key = "'category'", unless = "#result == null") public Map getCategoryStats() { EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program") .termsAggregation("category_count", "categoryId", 100); Map aggResult = esService.searchAggregation(queryBuilder); // 解析结果 Map stats = new HashMap<>(); AggregationResult categoryAgg = aggResult.get("category_count"); if (categoryAgg != null && categoryAgg.getBuckets() != null) { for (AggregationResult.BucketData bucket : categoryAgg.getBuckets()) { stats.put(bucket.getKey(), bucket.getDocCount()); } } return stats; } private List getHotKeywords() { // 从统计服务获取热门搜索词 return Arrays.asList("演唱会", "话剧", "音乐剧", "相声", "脱口秀"); } } ``` #### **监控与告警** ```java /** * ES 监控服务 */ @Slf4j @Service public class EsMonitorService { @Autowired private RestHighLevelClient client; @Autowired private MeterRegistry meterRegistry; /** * 1. 记录查询耗时指标 */ public List searchWithMetrics(EsQueryBuilder queryBuilder, Class clazz) { Timer.Sample sample = Timer.start(meterRegistry); try { // 执行查询 List result = esService.search(queryBuilder, clazz); // 记录成功指标 sample.stop(Timer.builder("es.query.duration") .tag("index", queryBuilder.getIndexName()) .tag("status", "success") .register(meterRegistry)); // 记录结果数量 meterRegistry.gauge("es.query.result.size", Tags.of("index", queryBuilder.getIndexName()), result.size()); return result; } catch (Exception e) { // 记录失败指标 sample.stop(Timer.builder("es.query.duration") .tag("index", queryBuilder.getIndexName()) .tag("status", "error") .register(meterRegistry)); meterRegistry.counter("es.query.error", Tags.of("index", queryBuilder.getIndexName())).increment(); throw e; } } /** * 2. 集群健康检查 */ @Scheduled(fixedRate = 60000) public void checkClusterHealth() { try { ClusterHealthRequest request = new ClusterHealthRequest(); ClusterHealthResponse response = client.cluster() .health(request, RequestOptions.DEFAULT); String status = response.getStatus().name(); int numberOfNodes = response.getNumberOfNodes(); int activeShards = response.getActiveShards(); // 记录指标 meterRegistry.gauge("es.cluster.nodes", numberOfNodes); meterRegistry.gauge("es.cluster.shards.active", activeShards); // 状态告警 if ("RED".equals(status)) { sendAlert("ES 集群状态异常: RED", String.format("节点数: %d, 活跃分片: %d", numberOfNodes, activeShards)); } else if ("YELLOW".equals(status)) { log.warn("ES 集群状态: YELLOW,可能有副本未分配"); } log.info("ES 集群健康检查: status={}, nodes={}, activeShards={}", status, numberOfNodes, activeShards); } catch (Exception e) { log.error("ES 集群健康检查失败", e); sendAlert("ES 集群健康检查失败", e.getMessage()); } } /** * 3. 索引状态监控 */ @Scheduled(fixedRate = 300000) public void checkIndexStats() { try { IndicesStatsRequest request = new IndicesStatsRequest(); IndicesStatsResponse response = client.indices() .stats(request, RequestOptions.DEFAULT); response.getIndices().forEach((indexName, stats) -> { long docCount = stats.getPrimaries().getDocs().getCount(); long storeSizeBytes = stats.getPrimaries().getStore().getSizeInBytes(); meterRegistry.gauge("es.index.docs.count", Tags.of("index", indexName), docCount); meterRegistry.gauge("es.index.store.size.bytes", Tags.of("index", indexName), storeSizeBytes); log.debug("索引 [{}] 统计: 文档数={}, 存储大小={}MB", indexName, docCount, storeSizeBytes / 1024 / 1024); }); } catch (Exception e) { log.error("索引状态监控失败", e); } } /** * 4. 慢查询日志收集 */ public void logSlowQuery(String indexName, String dsl, long elapsed) { if (elapsed > 3000) { // 超过3秒 log.warn("[慢查询] 索引: {}, 耗时: {}ms, DSL: {}", indexName, elapsed, dsl); // 发送告警 if (elapsed > 10000) { // 超过10秒 sendAlert("ES 慢查询告警", String.format("索引: %s, 耗时: %dms", indexName, elapsed)); } } } private void sendAlert(String title, String content) { // 发送告警通知(钉钉、企业微信、邮件等) log.error("【告警】{}: {}", title, content); } } ``` ## **十一、总结** ![[7-Blog/后端与微服务/assets/Elasticsearch_知识总结-3b295165.jpg]]